【发布时间】:2017-02-12 12:55:56
【问题描述】:
我需要创建一个从多个主题读取的 Spark 流,并为每个主题使用不同的解码器(每个主题包含不同的 avro 编码对象):
def decode_avro(message):
schem = avro.schema.parse(open("error_list.avsc").read())
bytes_reader = io.BytesIO(message)
decoder = avro.io.BinaryDecoder(bytes_reader)
reader = avro.io.DatumReader(schem)
return reader.read(decoder)
ssc = StreamingContext(sc, 2)
kvs = KafkaUtils.createDirectStream(ssc, [topic, topic2], {
"metadata.broker.list": brokers}, valueDecoder = decode_avro)
我不知道是否可以为每个主题指定不同的解码器回调,或者是否可以知道解码器函数上的主题名称(这样我可以将主题名称用于 avro 模式文件并在同一个函数中解码所有消息)
谢谢
【问题讨论】:
-
我也面临着同样的挑战。我看到这个问题已经超过 1 年了。你是如何绕过这个障碍的?
-
我们最终没有使用这种方法(甚至目前根本没有使用 Kafka)。我认为有一个 try/catch 系统会在引发异常时跳转到下一个解码器。是一个丑陋的解决方案,但我没有找到更好的解决方案!
-
好的,感谢您的更新。我找到了一个合适的解决方案,所以我会在这里添加它作为答案。
标签: python apache-spark apache-kafka avro