【问题标题】:Kafka Integration in Apache HeronApache Heron 中的 Kafka 集成
【发布时间】:2019-03-05 02:00:22
【问题描述】:

我正在尝试将 Kafka 与 Heron 拓扑集成。但是,我找不到最新版本的 Heron (0.17.5) 的任何示例。是否有任何可以共享的示例或有关如何实现自定义 Kafka Spout 和 Kafka Bolt 的任何建议?

编辑 1:

我相信 Heron 中有意弃用了 KafkaSpoutKafkaBolt 以让位于新的 Streamlet API。我目前正在查看是否可以使用 Streamlet API 构建 KafkaSourceKafkaSink。但是,当我尝试在 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


【解决方案1】:

我设法使用 Streamlet API for Heron 完成了这项工作。我在这里发布相同的内容。希望它可以帮助其他面临同样问题的人。

卡夫卡来源

public class KafkaSource implements Source {

    private String streamName;

    private Consumer<String, String> kafkaConsumer;
    private List<String> kafkaTopic;

    private static final Logger LOGGER = Logger.getLogger("KafkaSource");

    @Override
    public void setup(Context context) {

        this.streamName = context.getStreamName();

        kafkaTopic = Arrays.asList(KafkaProperties.KAFKA_TOPIC);

        Properties props = new Properties();
        props.put("bootstrap.servers", KafkaProperties.BOOTSTRAP_SERVERS);
        props.put("group.id", KafkaProperties.CONSUMER_GROUP_ID);
        props.put("enable.auto.commit", KafkaProperties.ENABLE_AUTO_COMMIT);
        props.put("auto.commit.interval.ms", KafkaProperties.AUTO_COMMIT_INTERVAL_MS);
        props.put("session.timeout.ms", KafkaProperties.SESSION_TIMEOUT);
        props.put("key.deserializer", KafkaProperties.KEY_DESERIALIZER);
        props.put("value.deserializer", KafkaProperties.VALUE_DESERIALIZER);
        props.put("auto.offset.reset", KafkaProperties.AUTO_OFFSET_RESET);
        props.put("max.poll.records", KafkaProperties.MAX_POLL_RECORDS);
        props.put("max.poll.interval.ms", KafkaProperties.MAX_POLL_INTERVAL_MS);

        this.kafkaConsumer = new KafkaConsumer<>(props);

        kafkaConsumer.subscribe(kafkaTopic);
    }

    @Override
    public Collection get() {

        List<String> kafkaRecords = new ArrayList<>();

        ConsumerRecords<String, String> records = kafkaConsumer.poll(Long.MAX_VALUE);

        for (ConsumerRecord<String, String> record : records) {
            String rVal = record.value();
            kafkaRecords.add(rVal);
        }

        return kafkaRecords;
    }

    @Override
    public void cleanup() {
        kafkaConsumer.wakeup();
    }
}

【讨论】:

  • 您知道如何使用 Java API 而不是 Streamlet API 将 KafkaSpout 集成到 Heron 吗?
猜你喜欢
  • 2017-01-23
  • 2023-04-05
  • 2017-04-02
  • 2021-01-11
  • 1970-01-01
  • 2017-12-12
  • 1970-01-01
  • 1970-01-01
  • 2017-06-07
相关资源
最近更新 更多