【发布时间】:2017-08-20 21:02:06
【问题描述】:
我正在尝试如下运行 flink 作业以从 Apache Kafka 读取数据并打印:
Java 程序
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
Properties properties = new Properties();
properties.setProperty("bootstrap.servers", "test.net:9092");
properties.setProperty("group.id", "flink_consumer");
properties.setProperty("zookeeper.connect", "dev.com:2181,dev2.com:2181,dev.com:2181/dev2");
properties.setProperty("topic", "topic_name");
DataStream<String> messageStream = env.addSource(new FlinkKafkaConsumer082<>("topic_name", new SimpleStringSchema(), properties));
messageStream.rebalance().map(new MapFunction<String, String>() {
private static final long serialVersionUID = -6867736771747690202L;
public String map(String value) throws Exception {
return "Kafka and Flink says: " + value;
}
}).print();
env.execute();
Scala 代码
var properties = new Properties();
properties.setProperty("bootstrap.servers", "msg01.staging.bigdata.sv2.247-inc.net:9092");
properties.setProperty("group.id", "flink_consumer");
properties.setProperty("zookeeper.connect", "host33.dev.swamp.sv2.tellme.com:2181,host37.dev.swamp.sv2.tellme.com:2181,host38.dev.swamp.sv2.tellme.com:2181/staging_sv2");
properties.setProperty("topic", "sv2.staging.rtdp.idm.events.omnichannel");
var env = StreamExecutionEnvironment.getExecutionEnvironment();
var stream:DataStream[(String)] = env
.addSource(new FlinkKafkaConsumer082[String]("sv2.staging.rtdp.idm.events.omnichannel", new SimpleStringSchema(), properties));
stream.print();
env.execute();
每当我在 Eclipse 中的应用程序中运行它时,我都会看到以下内容:
03/27/2017 20:06:19 作业执行切换到状态 RUNNING。
03/27/2017 20:06:19 来源:自定义来源 -> 接收器:未命名 (1/4) 切换到 SCHEDULED 2017 年 3 月 27 日 20:06:19 来源:自定义来源 -> 接收器:未命名(1/4)切换到部署 2017 年 3 月 27 日 20:06:19 来源:自定义来源 -> 接收器:未命名(2/4)切换到 SCHEDULED 2017 年 3 月 27 日 20:06:19 来源:自定义来源 -> 接收器:未命名(2/4)切换到部署 2017 年 3 月 27 日 20:06:19 来源:自定义来源 -> 接收器:未命名(3/4)切换到 SCHEDULED 2017 年 3 月 27 日 20:06:19 来源:自定义来源 -> 接收器:未命名(3/4)切换到部署 2017 年 3 月 27 日 20:06:19 来源:自定义来源 -> 接收器:未命名(4/4)切换到 SCHEDULED 2017 年 3 月 27 日 20:06:19 来源:自定义来源 -> 接收器:未命名(4/4)切换到部署 2017 年 3 月 27 日 20:06:19 来源:自定义来源 -> 接收器:未命名(4/4)切换到 RUNNING 2017 年 3 月 27 日 20:06:19 来源:自定义来源 -> 接收器:未命名(2/4)切换到 RUNNING 2017 年 3 月 27 日 20:06:19 来源:自定义来源 -> 接收器:未命名(1/4)切换到 RUNNING 2017 年 3 月 27 日 20:06:19 来源:自定义来源 -> 接收器:未命名(3/4)切换到 RUNNING
我的问题是:
1) 为什么我在所有情况下都看到 4 个接收器实例(计划、部署和运行)。
2) 对于 Apache Kafka 中收到的每一行,我看到这里被打印多次,大部分是 4 次。什么原因?
理想情况下,我只想读取每行一次并对其进行进一步处理。任何输入/帮助都将是可观的!
【问题讨论】:
-
你的主题有多少个分区?
-
我使用的主题有 PartitionCount:6 & ReplicationFactor:2
标签: apache-kafka apache-flink apache-kafka-connect bigdata