Spark读取Delta CDC写入Azure Event Hub缺失body属性报错如何解决?
解决方案
错误原因
Azure Event Hub 官方Spark连接器写入时强制要求输入DataFrame必须包含
body列,且列类型为二进制(binary),你从Delta变更数据捕获(CDC) feed 读取到的DataFrame只有Delta自带的CDC元字段(_change_type、_commit_version、_commit_timestamp)和原表业务字段,不存在符合要求的body列,因此触发校验异常。
修复步骤
- 首先导入需要的Spark SQL函数
from pyspark.sql.functions import to_json, struct, col
- 对读取到的CDC流数据做转换,生成符合要求的
body列:将整行需要传输的字段序列化为JSON字符串,再转换为二进制类型赋值给body列
# 转换得到适配EventHub写入要求的流DataFrame eh_ready_df = df.select( # *[col(c) for c in df.columns] 表示取原CDC流的所有字段,若只需部分字段可手动指定列名 to_json(struct(*[col(c) for c in df.columns])).cast("binary").alias("body") # 可选拓展:如需指定EventHub分区键可额外添加 partitionKey 列,如需添加自定义消息属性可添加 properties 列 )
- 用转换后的流DataFrame执行写入操作即可
eh_ready_df.writeStream \ .format("eventhubs") \ .option("checkpointLocation", checkpointLocation) \ .outputMode("append") \ .options(**ehConf) \ .start()
补充说明
如果不需要传输CDC元字段,仅需要原表的业务数据,可以手动指定struct内的字段,示例:
# 仅保留原表的id、order_amount、create_time三个业务字段写入EventHub eh_ready_df = df.select( to_json(struct("id", "order_amount", "create_time")).cast("binary").alias("body") )
消费端读取EventHub消息时,只需将body字段从二进制转换为字符串,再做JSON解析即可拿到完整的结构化数据。
内容的提问来源于stack exchange,提问作者Akhilesh Jaiswal
相关产品推荐
相关产品推荐

