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

如何在DBR 14.3LTS中复刻DBR 11.3LTS的Autoloader CSV读取行为

问题描述

从DBR 11.3LTS升级到14.3LTS后,Databricks Autoloader读取CSV文件出现列数超限相关问题:

  • 旧版本可自动跳过因含大量分号导致列数超标的行,且已有自定义逻辑将此类行分流至坏表
  • 新版本必须配置cloudFiles.maxColumns,但出现两种极端问题:
    • 设置过大值(如2147483647)时触发**"Requested array size exceeds VM limit"**错误
    • 设置较小值(如50000)时触发解析异常,错误信息翻译如下:

      由以下原因引起:com.univocity.parsers.common.TextParsingException: java.lang.ArrayIndexOutOfBoundsException - 20480
      提示:处理的列数可能已超过20480列的限制。请使用settings.setMaxColumns(int)定义输入允许的最大列数
      请确保你的配置正确,分隔符、引号和转义序列与你尝试解析的输入格式匹配

解决方案(无需大幅改动复刻旧版行为)

1. 合理设置列数上限+保留PERMISSIVE模式

不需要设置极端大值,先通过测试确认实际坏行的最大列数(比如设为10000),同时保持mode=PERMISSIVE,配合原有坏行分流逻辑:

  • 在Autoloader的ReadOptions中添加配置:
    "ReadOptions": {
      // ... 保留原有配置
      "mode": "PERMISSIVE",
      "cloudFiles.maxColumns": "10000",
      "columnNameOfCorruptRecord": "_corrupt_record"
    }
    
  • 继续使用原有自定义逻辑,将含_corrupt_record字段的行分流至坏表

2. 适配Univocity解析器底层配置

DBR14.3依赖的Univocity解析器对列数限制更严格,可直接配置解析器参数避免缓冲区溢出:

"ReadOptions": {
  // ... 保留原有配置
  "parser.columnLimit": "10000",
  "parser.maxCharsPerColumn": "1000000"
}

或者在SparkSession初始化时全局设置:

spark.conf.set("spark.sql.csv.parser.columnLimit", "10000")
spark.conf.set("spark.sql.csv.parser.maxCharsPerColumn", "1000000")

3. 临时回退到旧版解析器(过渡方案)

DBR14.x支持切换回旧版CSV解析器,直接复刻DBR11.3的解析行为,无需调整列数限制:

"ReadOptions": {
  // ... 保留原有配置
  "spark.sql.csv.useLegacyParser": "true"
}

注意:该配置属于过渡性选项,未来版本可能移除,仅建议短期使用。

调整后的配置示例
{"Name": "eventlog_n4_wlf_csv", 
 "DataProductName": "some_data_prod", 
 "LookbackRange": 21, 
 "Database": "some_db", 
 "DatasetPath": "some_path", 
 "CheckpointPath": "some_checkpoint", 
 "BadPath": "some_bad_path", 
 "BadTablePath": "some_bad_table_path", 
 "PartitionColumns": ["Year", "Month", "Day"], 
 "SchemaColumns": [{"colName": "Year", "nullable": false, "dataType": "int"}, {"colName": "Month", "nullable": false, "dataType": "int"}, {"colName": "Day", "nullable": false, "dataType": "int"}, {"colName": "SerialNumber", "nullable": false, "dataType": "int"}, {"colName": "MaterialNumber", "nullable": false, "dataType": "int"}, {"colName": "EventSourceHash", "nullable": false, "dataType": "string"}, {"colName": "Category", "nullable": true, "dataType": "int"}, {"colName": "CategoryString", "nullable": true, "dataType": "string"}, {"colName": "EventCode", "nullable": true, "dataType": "int"}, {"colName": "EventIdentifier", "nullable": true, "dataType": "bigint"}, {"colName": "EventType", "nullable": true, "dataType": "int"}, {"colName": "Logfile", "nullable": false, "dataType": "string"}, {"colName": "Message", "nullable": true, "dataType": "string"}, {"colName": "RecordNumber", "nullable": false, "dataType": "bigint"}, {"colName": "SourceName", "nullable": true, "dataType": "string"}, {"colName": "TimeGenerated", "nullable": false, "dataType": "timestamp"}, {"colName": "TimeZoneOffset", "nullable": true, "dataType": "int"}, {"colName": "Type", "nullable": true, "dataType": "string"}, {"colName": "User", "nullable": true, "dataType": "string"}, {"colName": "Data", "nullable": true, "dataType": "string"}, {"colName": "FileName", "nullable": false, "dataType": "string"}, {"colName": "FileUploadDate", "nullable": false, "dataType": "timestamp"}], 
 "Properties": {"delta.enableChangeDataFeed": "true", "delta.checkpoint.writeStatsAsStruct": "true", "delta.targetFileSize": "32mb", "delta.autoOptimize.autoCompact": "true"}, 
 "WriteOptions": {}, 
 "Sources": {"csv_source": {"Path": "/mnt/raw/some_data_product_path/exp/day=*/materialnum=*/serialnum=*/", 
                           "ReadOptions": {"useStrictGlobber": "true", 
                                           "header": "true", 
                                           "sep": ";", 
                                           "cloudFiles.partitionColumns": "day,materialnum,serialnum", 
                                           "cloudFiles.schemaEvolutionMode": "none", 
                                           "cloudFiles.schemaLocation": "/mnt/gold/tables/data_product/base_table/_schemaLocation", 
                                           "cloudFiles.useNotifications": "false", 
                                           "cloudFiles.maxFileAge": "28 days", 
                                           "cloudFiles.maxBytesPerTrigger": "8589934592", 
                                           "cloudFiles.maxFilesPerTrigger": "100", 
                                           "mode": "PERMISSIVE",
                                           "columnNameOfCorruptRecord": "_corrupt_record",
                                           "cloudFiles.maxColumns": "10000"}}}, 
 "Monitoring": {"TableName": "monitoring.processing_statistics", "TablePath": "/mnt/gold/tables/monitoring/processing_statistics"}}
调整后的测试代码示例
from pyspark.sql.types import StructType, StructField, IntegerType, StringType, LongType

# 定义schema,添加_corrupt_record字段标记坏行
schema = StructType(
    [
        StructField("Category", IntegerType(), True),
        StructField("CategoryString", StringType(), True),
        StructField("EventCode", IntegerType(), True),
        StructField("EventIdentifier", LongType(), True),
        StructField("EventType", IntegerType(), True),
        StructField("Logfile", StringType(), True),
        StructField("Message", StringType(), True),
        StructField("RecordNumber", LongType(), True),
        StructField("SourceName", StringType(), True),
        StructField("TimeGenerated", StringType(), True),
        StructField("Type", StringType(), True),
        StructField("User", StringType(), True),
        StructField("Data", StringType(), True),
        StructField("_corrupt_record", StringType(), True)
    ]
)

df = (
    spark.read.option("header", True)
    .option("delimiter", ";")
    .option("mode", "PERMISSIVE")
    .option("encoding", "UTF-8")
    .option("columnNameOfCorruptRecord", "_corrupt_record")
    .option("maxColumns", "10000")
    .csv("dbfs:/FileStore/asarkar/sample.csv", schema=schema)
)

# 复用原有逻辑分流好坏行
good_df = df.filter(df._corrupt_record.isNull())
bad_df = df.filter(df._corrupt_record.isNotNull())

display(good_df)

内容的提问来源于stack exchange,提问作者Arpan Sarkar

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 22:53:09