无法在Databricks查看Event Hub生成数据:字节数组异常及后续问题
解决Databricks读取Event Hubs模拟Yellow Taxi数据流的错误
第一个错误:Input byte array has wrong 4-byte ending unit
这个错误的核心原因是Event Hubs的Generate Data功能生成的消息体是Base64编码的二进制数据,直接通过cast("string")转换会触发Base64解码格式错误。正确的处理步骤是先对二进制body做Base64解码,再转为字符串。
修改后的代码如下:
from pyspark.sql.functions import unbase64, col connectionString = 'Endpoint=sb://carlsspace.servicebus.windows.net/;SharedAccessKeyName=myaccesspolicy;SharedAccessKey=lfZTXXXXXXXXXXXXXXXXXXXVaNv9wRM+AEhJ3XsFc=;EntityPath=carlshub' ehConf = {} # 适配2.3.15+版本的加密要求 ehConf['eventhubs.connectionString'] = sc._jvm.org.apache.spark.eventhubs.EventHubsUtils.encrypt(connectionString) df = spark \ .readStream \ .format("eventhubs") \ .options(**ehConf) \ .load() # 先解码Base64,再转换为字符串 df = df.withColumn("decoded_body", unbase64(col("body")).cast("string")) display(df.select("decoded_body"))
连接字符串加密的关键注意事项
- 严格匹配Spark Event Hubs Connector版本:2.3.15及以上版本必须使用
EventHubsUtils.encrypt()加密连接字符串;低于该版本直接传入明文连接串即可。 - Databricks社区版默认Spark版本通常较高,务必遵循高版本的加密要求,否则会触发授权或格式类错误。
额外排查方向
- 验证共享访问权限:确保
myaccesspolicy拥有Event Hub的Listen权限,无权限会导致连接失败。 - 核对EntityPath:必须与你的Event Hub名称完全一致,注意大小写敏感。
- 查看流监控日志:在Databricks中点击
display生成的流图表,进入监控界面查看完整错误栈,这能帮你定位后续出现的新问题。
内容的提问来源于stack exchange,提问作者Patterson
相关产品推荐
相关产品推荐

