ADF与Autoloader集成时如何添加可追溯列
问题根源
这不是Parquet格式的限制,核心原因是你在调用Auto Loader读取数据时,显式传入了仅包含源SQL Server表结构的StructType,Auto Loader默认只会读取该schema定义内的字段,ADF生成的三个额外列在读取阶段就被过滤掉了,仅会被存入你配置的_rescued_data救援列中,不会作为顶层字段出现在DataFrame里。而你设置的mergeSchema=true仅作用于Delta表写入阶段,只能将当前DataFrame中存在、但Delta表不存在的字段合并到表结构中,读取阶段就缺失的字段自然无法写入。
可用解决方案
你不需要将存储格式改为CSV或JSON,可根据场景选择以下方案解决:
- 方案1(最推荐,性能最优):手动扩展读取schema
在你当前使用的源表StructType中,新增三个ADF派生列的字段定义,与ADF输出的列名、类型保持一致即可,示例如下:
改完后Auto Loader会正常读取这三个字段,写入时from pyspark.sql.types import StructField, StringType, TimestampType # 原有源表schema定义保持不变,新增以下三个字段 schema = schema.add(StructField("pipeline_name", StringType(), True)) \ .add(StructField("runid", StringType(), True)) \ .add(StructField("trigger_time", TimestampType(), True))mergeSchema会自动将字段同步到Delta表结构中。 - 方案2(适配schema频繁变化场景):开启Auto Loader自动schema推断
去掉显式传入的固定schema,开启Auto Loader的字段推断和schema进化能力,配置如下:
该方案无需手动维护schema,但首次运行时需要扫描文件推断类型,大数据量下性能会有损耗,且存在字段类型推断不符合预期的风险。df = (spark .readStream .format("cloudFiles") .options(**cloudFile) .option("cloudFiles.inferColumnTypes", "true") .option("cloudFiles.schemaEvolutionMode", "addNewColumns") .option("rescueDataColumn", "_rescued_data") .load(sourceFilePath) ) - 方案3(无需修改读取配置,临时兼容):从救援列解析字段
你已经开启了_rescued_data配置,所有不在指定schema内的字段都会以JSON格式存储在该列中,可直接解析获取目标字段:df = ( df.withColumn("audit_fileName", input_file_name()) .withColumn("audit_createdTimestamp", current_timestamp()) # 新增以下三行从救援列解析ADF派生字段 .withColumn("pipeline_name", col("_rescued_data").getField("pipeline_name")) .withColumn("runid", col("_rescued_data").getField("runid")) .withColumn("trigger_time", col("_rescued_data").getField("trigger_time").cast(TimestampType())) )
内容的提问来源于stack exchange,提问作者Shrikant Kulkarni
相关产品推荐
相关产品推荐

