【发布时间】:2015-08-01 02:35:51
【问题描述】:
我在 python spark 应用程序中创建了一个 kafka 流,并且可以解析通过它的任何文本。
kafkaStream = KafkaUtils.createStream(ssc, zkQuorum, "spark-streaming-consumer", {topic: 1})
我想更改它以便能够解析来自 kafka 主题的 avro 消息。从文件中解析 avro 消息时,我会这样做:
reader = DataFileReader(open("customer.avro", "r"), DatumReader())
我是 python 和 spark 的新手,如何更改流以解析 avro 消息?另外,当从 Kafka 读取 Avro 消息时,如何指定要使用的模式???我以前在java中做过这一切,但是python让我很困惑。
编辑:
我尝试更改以包含 avro 解码器
kafkaStream = KafkaUtils.createStream(ssc, zkQuorum, "spark-streaming-consumer", {topic: 1},valueDecoder=avro.io.DatumReader(schema))
但我收到以下错误
TypeError: 'DatumReader' object is not callable
【问题讨论】:
-
您看到了什么错误?
标签: python apache-spark apache-kafka avro spark-streaming