【发布时间】:2019-03-05 02:00:22
【问题描述】:
我正在尝试将 Kafka 与 Heron 拓扑集成。但是,我找不到最新版本的 Heron (0.17.5) 的任何示例。是否有任何可以共享的示例或有关如何实现自定义 Kafka Spout 和 Kafka Bolt 的任何建议?
编辑 1:
我相信 Heron 中有意弃用了 KafkaSpout 和 KafkaBolt 以让位于新的 Streamlet API。我目前正在查看是否可以使用 Streamlet API 构建 KafkaSource 和 KafkaSink。但是,当我尝试在 Source 中创建 KafkaConsumer 时,出现以下异常。
Caused by: java.io.NotSerializableException: org.apache.kafka.clients.consumer.KafkaConsumer
at java.io.ObjectOutputStream.writeObject0(ObjectOutputStream.java:1184)
at java.io.ObjectOutputStream.defaultWriteFields(ObjectOutputStream.java:1548)
at java.io.ObjectOutputStream.writeSerialData(ObjectOutputStream.java:1509)
at java.io.ObjectOutputStream.writeOrdinaryObject(ObjectOutputStream.java:1432)
at java.io.ObjectOutputStream.writeObject0(ObjectOutputStream.java:1178)
at java.io.ObjectOutputStream.defaultWriteFields(ObjectOutputStream.java:1548)
at java.io.ObjectOutputStream.writeSerialData(ObjectOutputStream.java:1509)
at java.io.ObjectOutputStream.writeOrdinaryObject(ObjectOutputStream.java:1432)
at java.io.ObjectOutputStream.writeObject0(ObjectOutputStream.java:1178)
at java.io.ObjectOutputStream.writeObject(ObjectOutputStream.java:348)
at com.twitter.heron.api.utils.Utils.serialize(Utils.java:97)
编辑 2:
修复了上述问题。我在构造函数中初始化了KafkaConsumer,这是错误的。在setup() 方法中初始化它修复了它。
【问题讨论】:
-
如果 Heron 大致兼容 Storm,您有什么具体问题?
-
Storm 提供的 KafkaSpout 已弃用。
-
根据什么?文档是否提供了替代方案?
-
Streamlet API 似乎可以替代 Spout 和 Bolt。
标签: apache-kafka apache-storm kafka-consumer-api kafka-producer-api heron