【问题标题】:Convert avro serialized messages into json using python consumer使用python消费者将avro序列化消息转换为json
【发布时间】:2021-05-01 07:45:33
【问题描述】:
from kafka import KafkaConsumer
import json
import io

if __name__ == '__main__':

  # consumer = KafkaConsumer(
  #     'ldt_lm_mytable',
  #     bootstrap_servers = 'localhost:9092',
  #     auto_offset_reset = 'earliest',
  #     group_id = 'consumer_group_a')

  KAFKA_HOSTS = ['kafka:9092']
  KAFKA_VERSION = (0, 10)
  topic = "ldt_lm_mytable"

  consumer = KafkaConsumer(topic, bootstrap_servers=KAFKA_HOSTS, api_version=KAFKA_VERSION)

  for msg in consumer:
    print('Lead = {}'.format(json.loads(msg.value)))

没有打印。在将数据生成主题(Debezium)时,我正在使用 avro 转换器。我已经尝试了一些来自互联网的转换器。但这些都行不通。其中之一是这样的

bytes_reader = io.BytesIO(consumer)
decoder = avro.io.BinaryDecoder(bytes_reader)
reader = avro.io.DatumReader(schema)
decoded_data = reader.read(decoder)

在这个转换器中,我将从哪里获得“模式”变量的值?如何加载那个“avro”包?而那个'io.BytesIO'给了我一个像

这样的错误
Traceback (most recent call last):
File "consumer.py", line 19, in <module>
bytes_reader = io.BytesIO(consumer)
TypeError: a bytes-like object is required, not 'KafkaConsumer'

提前致谢!

【问题讨论】:

  • 如果这是您想要使用的格式,为什么不直接使用 json 转换器?如果数据未使用 Confluent Avro 格式进行序列化,您正在寻找 io.BytesIO(msg.value)

标签: python apache-kafka kafka-consumer-api avro kafka-python


【解决方案1】:

假设 Debezium 连接器在 Kafka Connect 中使用标准 io.confluent.connect.avro.AvroConverter,那么您需要使用与 Confluent Schema Registry 一起使用的 Avro 反序列化器。

这是消费者from here的示例:

from confluent_kafka.avro import AvroConsumer
from confluent_kafka.avro.serializer import SerializerError


c = AvroConsumer({
    'bootstrap.servers': 'mybroker,mybroker2',
    'group.id': 'groupid',
    'schema.registry.url': 'http://127.0.0.1:8081'})

c.subscribe(['my_topic'])

while True:
    try:
        msg = c.poll(10)

    except SerializerError as e:
        print("Message deserialization failed for {}: {}".format(msg, e))
        break

    if msg is None:
        continue

    if msg.error():
        print("AvroConsumer error: {}".format(msg.error()))
        continue

    print(msg.value())

c.close()

【讨论】:

    猜你喜欢
    • 2023-01-10
    • 2019-08-14
    • 2019-08-30
    • 1970-01-01
    • 2011-02-17
    • 2020-11-22
    • 2016-09-06
    • 2019-05-11
    • 2019-07-12
    相关资源
    最近更新 更多