【问题标题】:Apache Flink - multiple output linesApache Flink - 多条输出线
【发布时间】: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


【解决方案1】:

如果您在LocalStreamEnvironment(在IDE 中调用StreamExecutionEnvironment.getExecutionEnvironment() 时得到)中运行程序,则所有运算符的默认并行度等于CPU 内核的数量。

因此,在您的示例中,每个运算符都并行化为四个子任务。在日志中,您会看到这四个子任务中的每一个的消息(3/4 表示这是总共四个任务中的第三个)。

您可以通过在每个单独的操作员上调用StreamExecutionEnvironment.setParallelism(int) 或调用setParallelism(int) 来控制子任务的数量。

鉴于您的程序,不应复制 Kafka 记录。每条记录只能打印一次。但是,由于记录是并行写入的,因此输出行以x&gt; 为前缀,其中x 表示发出该行的并行子任务的ID。

【讨论】:

  • 感谢您的意见!根据您的评论,我尝试了以下操作: 1) StreamExecutionEnvironment.setDefaultLocalParallelism(1);因为我只想打印或处理每条记录一次。这比我看到某些行只打印一次之前效果更好,但过了一段时间它再次开始表现相似并打印多行。 2) env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(1);它仍然多次打印输出行。最终我需要处理这个输出并将数据放入 cassandra,所以我主要关心的是它会处理和存储多行。请分享您的意见!
  • 您发布的不复制数据的程序。您确定数据没有在 Kafka 主题内复制吗?
  • 已检查,似乎 Kafka 中打印的事件也存在一些问题。非常感谢费边!
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-08-13
相关资源
最近更新 更多