如何在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
相关产品推荐
相关产品推荐

