【问题标题】:Kafka Storm Integration卡夫卡风暴集成
【发布时间】:2016-01-30 21:41:59
【问题描述】:

我正在尝试整合 Kafka 和 Storm。

我创建了一个 spout 来读取 Kafka 消息并将它们作为元组发出。

在运行风暴拓扑时,我收到以下异常

Caused by: java.io.NotSerializableException: kafka.javaapi.consumer.ZookeeperConsumerConnector
    at java.io.ObjectOutputStream.writeObject0(ObjectOutputStream.java:1184) ~[?:1.8.0_45]
    at java.io.ObjectOutputStream.defaultWriteFields(ObjectOutputStream.java:1548) ~[?:1.8.0_45]
    at java.io.ObjectOutputStream.writeSerialData(ObjectOutputStream.java:1509) ~[?:1.8.0_45]
    at java.io.ObjectOutputStream.writeOrdinaryObject(ObjectOutputStream.java:1432) ~[?:1.8.0_45]
    at java.io.ObjectOutputStream.writeObject0(ObjectOutputStream.java:1178) ~[?:1.8.0_45]
    at java.io.ObjectOutputStream.writeObject(ObjectOutputStream.java:348) ~[?:1.8.0_45]
    at backtype.storm.utils.Utils.javaSerialize(Utils.java:87) ~[storm-core-0.10.0.jar:0.10.0]
    ... 2 more

【问题讨论】:

    标签: java apache-kafka apache-storm


    【解决方案1】:

    您应该将 spout 中的不可序列化字段声明为瞬态。这些字段应该在你的 spout 的 open 方法中初始化。

    【讨论】:

    • 您好,感谢您的回复。有效。但是现在我在 java.util.concurrent.locks.AbstractQueuedSynchronizer$ConditionObject.reportInterruptAfterWait(AbstractQueuedSynchronizer.java:2014) ~[?:1.8.0_45] at java.util.concurrent.locks 处得到以下异常 java.lang.InterruptedException .AbstractQueuedSynchronizer$ConditionObject.await(AbstractQueuedSynchronizer.java:2048) ~[?:1.8.0_45] 这个异常发生在我在 Spout 的 nextTuple() 方法中使用的迭代器中。
    • 您可能还会遇到其他异常。我没有你的代码,所以我不能添加更多。
    猜你喜欢
    • 2017-05-14
    • 2014-03-15
    • 2015-06-18
    • 2017-07-07
    • 1970-01-01
    • 1970-01-01
    • 2019-04-11
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多