【问题标题】:spark streaming different value decoder per kafka topic每个 kafka 主题的火花流不同值解码器
【发布时间】: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


【解决方案1】:

我们也有这种情况,我们从不同的主题读取不同的消息格式,然后处理每个主题并将输出存储到每个源主题的专用存储中。 去这里的正确方法是创建多个流。在同一个应用程序中使用相同的 Spark 上下文按主题流式传输。 每个流都会获得相关的 ValueDecoder,如果它们共享相同的格式,您仍然可以从多个主题中读取。

【讨论】:

  • 感谢您的回复。我会投票赞成它,当我有时间时,我会测试它。
猜你喜欢
  • 2019-10-11
  • 2015-10-13
  • 2017-04-27
  • 1970-01-01
  • 2020-04-11
  • 2016-08-21
  • 2023-03-18
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多