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

在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 01:31:21