【问题标题】:How to print kafka topic data in Flink program?如何在 Flink 程序中打印 kafka 主题数据?
【发布时间】:2018-11-05 11:27:12
【问题描述】:

我通过这条指令创建了一个主题:

C:\kafka_2.12-0.10.2.1>.\bin\windows\kafka-console-producer.bat --broker-list localhost:9092 --topic test < C:\User11\Desktop\Data.csv

然后我测试了该主题是否具有正确的数据。之后,我想在Flink程序中打印主题。我的程序是这样的:

 try{
    StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

    Properties properties = new Properties();
    DataStream<String> stream = env
            .addSource(new FlinkKafkaConsumer09<String>("test", new SimpleStringSchema(),properties));

           stream.print();
    env.execute();
    } catch (Exception e) {
        e.printStackTrace();
    }

但是我得到了这个INFO(因为INFO太长我不得不写一些):

[main] INFO org.apache.flink.streaming.api.environment.LocalStreamEnvironment - 在本地嵌入式 Flink 迷你集群上运行作业 [main] INFO org.apache.flink.runtime.minicluster.MiniCluster - 启动 Flink 迷你集群 [main] INFO org.apache.flink.runtime.minicluster.MiniCluster - 启动 Metrics Registry [main] INFO org.apache.flink.runtime.metrics.MetricRegistryImpl - 没有配置指标报告器,不会公开/报告任何指标。 [main] INFO org.apache.flink.runtime.minicluster.MiniCluster - 启动 RPC 服务 [flink-akka.actor.default-dispatcher-2] 信息 akka.event.slf4j.Slf4jLogger - Slf4jLogger 已启动 [main] INFO org.apache.flink.runtime.minicluster.MiniCluster - 启动高可用服务 [main] INFO org.apache.flink.runtime.blob.BlobServer - 创建 BLOB 服务器存储目录 C:\Users\user11\AppData\Local\Temp\blobStore-a02ff126-35cc-4c1b-b300-8689d19ff5d2 [main] INFO org.apache.flink.runtime.blob.BlobServer - 在 0.0.0.0:57907 启动 BLOB 服务器 - 最大并发请求数:50 - 最大积压:1000

另外,我也看到了这个链接,但它并没有解决我的问题: How to access/read kafka topic data from flink?

你能告诉我这里有什么问题吗?

谢谢。

【问题讨论】:

  • 看起来它对我有用...您可能需要告诉 Flink 从主题的开头阅读ci.apache.org/projects/flink/flink-docs-stable/dev/connectors/…
  • 谢谢@cricket_007,通过添加这一行,我可以打印主题的全部内容:“myconsumer.setStartFromEarliest();”如何逐行访问主题内容?
  • Kafka 服务器意味着从 Kafka 网站下载最新(或只是更新)版本。你可以阅读 Flink 文档,它告诉你在 Maven 中导入什么来处理特定的 Kafka 版本
  • flink-connector-kafka-0.11_2.11 仍然可以与较新的 Kafka 服务器一起使用

标签: java apache-kafka apache-flink


【解决方案1】:

问题解决了。首先,我用这个命令填充了 Kafka 主题:

/home/kafka_2.11-2.0.0/bin/kafka-console-producer.sh --broker-list 10.32.0.2:9092,10.32.0.3:9092,10.32.0.4:9092 --topic flinkTopic < transactions2.csv

然后,使用此代码,我可以打印 Kafka 主题:

 final StreamExecutionEnvironment env = 
 StreamExecutionEnvironment.getExecutionEnvironment();
 Properties prop = new Properties();
 prop.setProperty("bootstrap.servers", 
 "10.32.0.2:9092,10.32.0.3:9092,10.32.0.4:9092");
 prop.setProperty("group.id", "test");
    FlinkKafkaConsumer<String> myConsumer= new FlinkKafkaConsumer<> 
  ("flinkTopic", new SimpleStringSchema(),prop);
    myConsumer.setStartFromEarliest();
    DataStream<String> stream = env.addSource(myConsumer);
    stream.print();
    env.execute("Flink Streaming Java API Skeleton");

我希望它对其他人有用。

【讨论】:

    猜你喜欢
    • 2021-07-30
    • 2022-06-30
    • 1970-01-01
    • 1970-01-01
    • 2020-09-18
    • 1970-01-01
    • 2022-01-26
    • 1970-01-01
    • 2019-12-07
    相关资源
    最近更新 更多