【问题标题】:Spark Python Avro Kafka DeserialiserSpark Python Avro Kafka 反序列化器
【发布时间】: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


【解决方案1】:

我遇到了同样的挑战 - 在 pyspark 中反序列化来自 Kafka 的 avro 消息,并使用 Confluent Schema Registry 模块的 Messageserializer 方法解决了这个问题,因为在我们的例子中,架构存储在 Confluent Schema Registry 中。

您可以在https://github.com/verisign/python-confluent-schemaregistry找到该模块

from confluent.schemaregistry.client import CachedSchemaRegistryClient
from confluent.schemaregistry.serializers import MessageSerializer
schema_registry_client = CachedSchemaRegistryClient(url='http://xx.xxx.xxx:8081')
serializer = MessageSerializer(schema_registry_client)


# simple decode to replace Kafka-streaming's built-in decode decoding UTF8 ()
def decoder(s):
    decoded_message = serializer.decode_message(s)
    return decoded_message

kvs = KafkaUtils.createDirectStream(ssc, ["mytopic"], {"metadata.broker.list": "xxxxx:9092,yyyyy:9092"}, valueDecoder=decoder)

lines = kvs.map(lambda x: x[1])
lines.pprint()

显然,正如您所看到的,这段代码使用的是新的、直接的方法,没有接收器,因此是 createdDirectStream(在https://spark.apache.org/docs/1.5.1/streaming-kafka-integration.html 上查看更多信息)

【讨论】:

  • 你说的那个库现在有点老了,好像没有维护了。
  • 哈哈,是的,我在 2.5 年前提到过那个图书馆,当时它还很“新鲜”。 :-)
【解决方案2】:

正如@Zoltan Fedor 在评论中提到的那样,提供的答案现在有点老了,因为它已经过去了 2.5 年。 confluent-kafka-python 库已经发展为原生支持相同的功能。给定代码中唯一需要的更改如下。

from confluent_kafka.avro.cached_schema_registry_client import CachedSchemaRegistryClient
from confluent_kafka.avro.serializer.message_serializer import MessageSerializer

然后,你可以改变这一行 -

kvs = KafkaUtils.createDirectStream(ssc, ["mytopic"], {"metadata.broker.list": "xxxxx:9092,yyyyy:9092"}, valueDecoder=serializer.decode_message)

我已经对其进行了测试,并且效果很好。我正在为将来可能需要它的任何人添加此答案。

【讨论】:

    【解决方案3】:

    如果您不考虑使用 Confluent Schema Registry 并且在文本文件或 dict 对象中有模式,您可以使用fastavro python 包来解码您的 Kafka 流的 Avro 消息:

    from pyspark.streaming.kafka import KafkaUtils
    from pyspark.streaming import StreamingContext
    import io
    import fastavro
    
    def decoder(msg):
        # here should be your schema
        schema = {
          "namespace": "...",
          "type": "...",
          "name": "...",
          "fields": [
            {
              "name": "...",
              "type": "..."
            },
          ...}
        bytes_io = io.BytesIO(msg)
        bytes_io.seek(0)
        msg_decoded = fastavro.schemaless_reader(bytes_io, schema)
        return msg_decoded
    
    session = SparkSession.builder \
                          .appName("Kafka Spark Streaming Avro example") \
                          .getOrCreate()
    
    streaming_context = StreamingContext(sparkContext=session.sparkContext,
                                         batchDuration=5)
    
    kafka_stream = KafkaUtils.createDirectStream(ssc=streaming_context,
                                                 topics=['your_topic_1', 'your_topic_2'],
                                                 kafkaParams={"metadata.broker.list": "your_kafka_broker_1,your_kafka_broker_2"},
                                                 valueDecoder=decoder)
    

    【讨论】:

      猜你喜欢
      • 2020-03-04
      • 2019-11-18
      • 2020-10-05
      • 1970-01-01
      • 2019-07-30
      • 2017-12-14
      • 2019-08-30
      • 2018-08-11
      • 1970-01-01
      相关资源
      最近更新 更多