Python 使用 avro 存储库反序列化 kafka 消息
Python deserialize kafka message with avro repository
我需要使用存储在存储库中的 avro 读取 Kafka 消息。
使用 kafka-python 2.0.2,我可以连接到 Kafka 主题并阅读消息,但我不知道如何解码它们。
from kafka import KafkaConsumer
consumer = KafkaConsumer('SOME-TOPIC',
other connection parameters,
auto_offset_reset= 'earliest')
# value_deserializer=lambda m: json.loads(m.decode('utf-8')))
# value_deserializer=lambda m: decode(m))
for msg in consumer:
print (msg)
我应该使用什么库?融合-kafka 1.5.0,avro-python3 1.10.1
如何进行?
- 识别消息的版本
- 连接到 avro 存储库
- 获取正确版本的 avro
- 用它来解码消息
这似乎有很多事情要做,有没有更简单的方法呢?
我很高兴能得到一个例子来指导我。
要连接到 avro 存储库,我有这些参数
- basic.auth.credentials.source
- schema.registry.basic.auth.user.info
- schema.registry.url
我使用了以下库:
pip install confluent-avro
pip install kafka-python
代码:
from kafka import KafkaConsumer
from confluent_avro import AvroKeyValueSerde, SchemaRegistry
from confluent_avro.schema_registry import HTTPBasicAuth
KAFKA_TOPIC = "SOME-TOPIC"
registry_client = SchemaRegistry(
"https://...",
HTTPBasicAuth("USER", "PASSWORD"),
headers={"Content-Type": "application/vnd.schemaregistry.v1+json"},
)
avroSerde = AvroKeyValueSerde(registry_client, KAFKA_TOPIC)
consumer = KafkaConsumer(KAFKA_TOPIC,
other connection parameters,
auto_offset_reset= 'earliest')
for msg in consumer:
v = avroSerde.value.deserialize(msg.value)
k = avroSerde.key.deserialize(msg.key)
print(msg.offset, msg.partition, k, v)
break
我需要使用存储在存储库中的 avro 读取 Kafka 消息。 使用 kafka-python 2.0.2,我可以连接到 Kafka 主题并阅读消息,但我不知道如何解码它们。
from kafka import KafkaConsumer
consumer = KafkaConsumer('SOME-TOPIC',
other connection parameters,
auto_offset_reset= 'earliest')
# value_deserializer=lambda m: json.loads(m.decode('utf-8')))
# value_deserializer=lambda m: decode(m))
for msg in consumer:
print (msg)
我应该使用什么库?融合-kafka 1.5.0,avro-python3 1.10.1 如何进行?
- 识别消息的版本
- 连接到 avro 存储库
- 获取正确版本的 avro
- 用它来解码消息
这似乎有很多事情要做,有没有更简单的方法呢? 我很高兴能得到一个例子来指导我。
要连接到 avro 存储库,我有这些参数
- basic.auth.credentials.source
- schema.registry.basic.auth.user.info
- schema.registry.url
我使用了以下库:
pip install confluent-avro
pip install kafka-python
代码:
from kafka import KafkaConsumer
from confluent_avro import AvroKeyValueSerde, SchemaRegistry
from confluent_avro.schema_registry import HTTPBasicAuth
KAFKA_TOPIC = "SOME-TOPIC"
registry_client = SchemaRegistry(
"https://...",
HTTPBasicAuth("USER", "PASSWORD"),
headers={"Content-Type": "application/vnd.schemaregistry.v1+json"},
)
avroSerde = AvroKeyValueSerde(registry_client, KAFKA_TOPIC)
consumer = KafkaConsumer(KAFKA_TOPIC,
other connection parameters,
auto_offset_reset= 'earliest')
for msg in consumer:
v = avroSerde.value.deserialize(msg.value)
k = avroSerde.key.deserialize(msg.key)
print(msg.offset, msg.partition, k, v)
break