PySpark无法反序列化Event Hub中Azure AvroEncoder编码数据问题
问题场景
我们通过独立Python作业将经azure.schemaregistry.encoder.avroencoder编码的Avro数据发送至Event Hub,使用同套解码器的另一独立Python消费者可正常完成反序列化,该场景下Avro编码器已正确配置Schema Registry。
独立生产者测试代码
import os from azure.eventhub import EventHubProducerClient, EventData from azure.schemaregistry import SchemaRegistryClient from azure.schemaregistry.encoder.avroencoder import AvroEncoder from azure.identity import DefaultAzureCredential os.environ["AZURE_CLIENT_ID"] = '' os.environ["AZURE_TENANT_ID"] = '' os.environ["AZURE_CLIENT_SECRET"] = '' token_credential = DefaultAzureCredential() fully_qualified_namespace = "" group_name = "testSchemaReg" eventhub_connection_str = "" eventhub_name = "" definition = """ {"namespace": "example.avro", "type": "record", "name": "User", "fields": [ {"name": "name", "type": "string"}, {"name": "favorite_number", "type": ["int", "null"]}, {"name": "favorite_color", "type": ["string", "null"]} ] }""" schema_registry_client = SchemaRegistryClient(fully_qualified_namespace, token_credential) avro_encoder = AvroEncoder(client=schema_registry_client, group_name=group_name, auto_register=True) eventhub_producer = EventHubProducerClient.from_connection_string( conn_str=eventhub_connection_str, eventhub_name=eventhub_name ) with eventhub_producer, avro_encoder: event_data_batch = eventhub_producer.create_batch() dict_content = {"name": "Bob", "favorite_number": 7, "favorite_color": "red"} event_data = avro_encoder.encode(dict_content, schema=definition, message_type=EventData) event_data_batch.add(event_data) eventhub_producer.send_batch(event_data_batch)
可正常反序列化的独立消费者代码
async def on_event(partition_context, event): print("Received the event: \"{}\" from the partition with ID: \"{}\"".format(event.body_as_str(encoding='UTF-8'), partition_context.partition_id)) print("message type is :") print(type(event)) dec = avro_encoder.decode(event) print("decoded msg:\n") print(dec) await partition_context.update_checkpoint(event) async def main(): client = EventHubConsumerClient.from_connection_string( "connection str" "topic name", consumer_group="$Default", eventhub_name="") async with client: await client.receive(on_event=on_event, starting_position="-1")
替换为Synapse Notebook PySpark消费者时遇到的问题
将独立Python消费者替换为运行在synapse notebook上的py-spark consumer后,出现以下问题:
- Spark内置的from_avro函数无法反序列化经Azure编码器编码的Avro消息。
- 尝试编写调用Azure编码器的UDF作为替代方案时,发现Azure编码器要求输入为EventData类型,而Spark通过Event Hub API读取到的数据为字节数组,无法直接传入解码,对应UDF代码如下:
@udf def decode(row_msg): encoder = AvroEncoder(client=schema_registry_client) encoder.decode(bytes(row_msg))
- 目前未查询到适用于Spark或其他分布式系统的反序列化组件正式文档,所有官方示例均面向独立客户端,需确认是否存在可用于Spark/Flink的适配连接器。
内容的提问来源于stack exchange,提问作者Kiran Gali
相关产品推荐
相关产品推荐

