You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.28 23:33:26