在Databricks中用Auto Loader(BinaryFile)直接写入Delta格式的问题
在Databricks中用Auto Loader处理Proto文件并直接写入Delta的问题
我在Databricks中使用Auto Loader的BinaryFile选项解码基于.proto的文件,目前已能通过foreach()和pandas库完成解码并写入CSV格式,但在直接写入Delta格式时遇到瓶颈。最终目标是直接写入Delta格式,完全跳过CSV中转步骤。
我尝试了几种方法,但均存在问题:
- 将pandas DataFrame转换为Spark DataFrame:需要依赖sparkContext创建DataFrame,但无法将sparkContext广播到工作节点,导致操作失败。
- 不使用pandas DataFrame:仍需在
foreach()内创建DataFrame,但加载操作分布在多个工作节点上,这种方式不可行。 - 使用UDF:解码后展开返回的字符串,但由于处理的是Spark非原生的proto格式,该方法不适用。
我也查阅了相关资源,但未找到针对foreach()和BinaryFile场景的有效方案:
- Delta-rs的Python版本尚未稳定,无法用于生产环境。
- PySpark pandas的
to_delta方法依然会触发上述第二种方法的问题。
若有相关解决方向,将不胜感激。
参考代码片段
cloudfile_options = { "cloudFiles.subscriptionId": subscription_ID, "cloudFiles.connectionString": queue_connection_string, "cloudFiles.format": "BinaryFile", "cloudFiles.tenantId":tenant_ID, "cloudFiles.clientId":client_ID, "cloudFiles.clientSecret":client_secret, "cloudFiles.resourceGroup": storage_resource_group, "cloudFiles.useNotifications" :"true" } reader_df = spark.readStream.format("cloudFiles") \ .options(**cloudfile_options) \ .load("some_storage_input_path") def decode_proto(self, row): with open(row['path'], 'rb') as f: # 执行解码操作 # 将解码后的字符串转为Json并通过pandas DataFrame写入存储 write_stream = reader_df.select("path") \ .writeStream \ .foreach(decode_proto) \ .option("checkpointLocation", checkpoint_path) \ .trigger(once=True) \ .start()
内容的提问来源于stack exchange,提问作者pavan
相关产品推荐
相关产品推荐

