【发布时间】: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