Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
caopu16
whatsmars
提交
e5909599
W
whatsmars
项目概览
caopu16
/
whatsmars
与 Fork 源项目一致
Fork自
武汉红喜 / whatsmars
通知
1
Star
0
Fork
0
代码
文件
提交
分支
Tags
贡献者
分支图
Diff
Issue
0
列表
看板
标记
里程碑
合并请求
0
Wiki
0
Wiki
分析
仓库
DevOps
项目成员
Pages
W
whatsmars
项目概览
项目概览
详情
发布
仓库
仓库
文件
提交
分支
标签
贡献者
分支图
比较
Issue
0
Issue
0
列表
看板
标记
里程碑
合并请求
0
合并请求
0
Pages
分析
分析
仓库分析
DevOps
Wiki
0
Wiki
成员
成员
收起侧边栏
关闭侧边栏
动态
分支图
创建新Issue
提交
Issue看板
前往新版Gitcode,体验更适合开发者的 AI 搜索 >>
提交
e5909599
编写于
9月 11, 2017
作者:
S
shenhongxi
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
kafka
上级
3595c061
变更
3
隐藏空白更改
内联
并排
Showing
3 changed file
with
104 addition
and
0 deletion
+104
-0
whatsmars-mq/pom.xml
whatsmars-mq/pom.xml
+27
-0
whatsmars-mq/src/main/java/com/itlong/whatsmars/mq/kafka/KafkaConsumer.java
...ain/java/com/itlong/whatsmars/mq/kafka/KafkaConsumer.java
+47
-0
whatsmars-mq/src/main/java/com/itlong/whatsmars/mq/kafka/KafkaProducer.java
...ain/java/com/itlong/whatsmars/mq/kafka/KafkaProducer.java
+30
-0
未找到文件。
whatsmars-mq/pom.xml
浏览文件 @
e5909599
...
...
@@ -65,6 +65,33 @@
<version>
0.4.2
</version>
</dependency>
<dependency>
<groupId>
org.apache.kafka
</groupId>
<artifactId>
kafka_2.8.2
</artifactId>
<version>
0.8.0
</version>
<exclusions>
<exclusion>
<groupId>
log4j
</groupId>
<artifactId>
log4j
</artifactId>
</exclusion>
</exclusions>
</dependency>
<dependency>
<groupId>
org.scala-lang
</groupId>
<artifactId>
scala-library
</artifactId>
<version>
2.8.2
</version>
</dependency>
<dependency>
<groupId>
com.yammer.metrics
</groupId>
<artifactId>
metrics-core
</artifactId>
<version>
2.2.0
</version>
</dependency>
<dependency>
<groupId>
com.101tec
</groupId>
<artifactId>
zkclient
</artifactId>
<version>
0.3
</version>
</dependency>
</dependencies>
...
...
whatsmars-mq/src/main/java/com/itlong/whatsmars/mq/kafka/KafkaConsumer.java
0 → 100644
浏览文件 @
e5909599
package
com.itlong.whatsmars.mq.kafka
;
import
kafka.consumer.ConsumerConfig
;
import
kafka.consumer.ConsumerIterator
;
import
kafka.consumer.KafkaStream
;
import
kafka.javaapi.consumer.ConsumerConnector
;
import
kafka.serializer.StringDecoder
;
import
kafka.utils.VerifiableProperties
;
import
java.util.HashMap
;
import
java.util.List
;
import
java.util.Map
;
import
java.util.Properties
;
public
class
KafkaConsumer
{
public
static
void
main
(
String
[]
args
)
{
Properties
props
=
new
Properties
();
props
.
put
(
"zookeeper.connect"
,
"127.0.0.1:2181"
);
//group 代表一个消费组
props
.
put
(
"group.id"
,
"jd-group"
);
props
.
put
(
"zookeeper.session.timeout.ms"
,
"4000"
);
props
.
put
(
"zookeeper.sync.time.ms"
,
"200"
);
props
.
put
(
"auto.commit.interval.ms"
,
"1000"
);
props
.
put
(
"auto.offset.reset"
,
"smallest"
);
props
.
put
(
"serializer.class"
,
"kafka.serializer.StringEncoder"
);
ConsumerConfig
config
=
new
ConsumerConfig
(
props
);
ConsumerConnector
consumer
=
kafka
.
consumer
.
Consumer
.
createJavaConsumerConnector
(
config
);
Map
<
String
,
Integer
>
topicCountMap
=
new
HashMap
<
String
,
Integer
>();
topicCountMap
.
put
(
"TestTopic"
,
1
);
StringDecoder
keyDecoder
=
new
StringDecoder
(
new
VerifiableProperties
());
StringDecoder
valueDecoder
=
new
StringDecoder
(
new
VerifiableProperties
());
Map
<
String
,
List
<
KafkaStream
<
String
,
String
>>>
consumerMap
=
consumer
.
createMessageStreams
(
topicCountMap
,
keyDecoder
,
valueDecoder
);
KafkaStream
<
String
,
String
>
stream
=
consumerMap
.
get
(
"TestTopic"
).
get
(
0
);
ConsumerIterator
<
String
,
String
>
it
=
stream
.
iterator
();
while
(
it
.
hasNext
())
{
System
.
out
.
println
(
it
.
next
().
message
());
}
System
.
out
.
println
(
"finished"
);
}
}
\ No newline at end of file
whatsmars-mq/src/main/java/com/itlong/whatsmars/mq/kafka/KafkaProducer.java
0 → 100644
浏览文件 @
e5909599
package
com.itlong.whatsmars.mq.kafka
;
import
kafka.javaapi.producer.Producer
;
import
kafka.producer.KeyedMessage
;
import
kafka.producer.ProducerConfig
;
import
java.util.Properties
;
public
class
KafkaProducer
{
public
static
void
main
(
String
[]
args
)
{
Properties
props
=
new
Properties
();
props
.
put
(
"metadata.broker.list"
,
"127.0.0.1:9092"
);
props
.
put
(
"serializer.class"
,
"kafka.serializer.StringEncoder"
);
props
.
put
(
"key.serializer.class"
,
"kafka.serializer.StringEncoder"
);
props
.
put
(
"request.required.acks"
,
"-1"
);
Producer
<
String
,
String
>
producer
=
new
Producer
<
String
,
String
>(
new
ProducerConfig
(
props
));
int
messageNo
=
100
;
final
int
COUNT
=
1000
;
while
(
messageNo
<
COUNT
)
{
String
key
=
String
.
valueOf
(
messageNo
);
String
data
=
"hello kafka message "
+
key
;
producer
.
send
(
new
KeyedMessage
<
String
,
String
>(
"TestTopic"
,
key
,
data
));
System
.
out
.
println
(
data
);
messageNo
++;
}
}
}
\ No newline at end of file
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录