Spark Structured Streaming读取CSV文件时inferSchema=true不生效报错问题
问题根因
Spark Structured Streaming 读取文件类数据源时默认关闭自动Schema推断能力,属于框架内置的设计限制,和参数写法、CSV文件格式无关。
流式场景下数据会持续流入,Spark为了避免后续流入文件Schema突变引发运行时错误,默认要求必须显式指定Schema,仅配置option("inferSchema", "true")不会生效。
解决方案
两种方案可选,第二种更适合生产环境使用:
方案1:开启全局流式Schema推断配置
开启spark.sql.streaming.schemaInference配置后,自动推断参数即可生效,Spark会扫描目标路径下所有已存在的CSV文件完成Schema推断。
代码示例:// 开启全局配置,必须在读取流之前设置 spark2.conf.set("spark.sql.streaming.schemaInference", "true") val streamDf = spark2 .readStream .format("csv") .option("header", "true") .option("delimiter", ",") .option("maxFilesPerTrigger", 1) .csv(path)注意:该方案启动时需要扫描路径下所有已有文件,文件量大时启动速度慢,且需要保证后续流入的CSV文件Schema和推断结果完全一致,否则会出现解析错误。
方案2:先静态读取推断Schema再传入流(官方推荐)
按照报错提示的方案,先通过批量静态读取的方式推断得到Schema,再将Schema传入流读取逻辑,兼顾自动推断的便捷性和运行稳定性。
代码示例:// 静态读取样本文件推断Schema,仅需扫描少量样本即可,速度快 val sampleCsvDf = spark2.read .format("csv") .option("header", "true") .option("inferSchema", "true") .csv(path) val csvSchema = sampleCsvDf.schema // 流式读取时传入提前推断好的Schema val streamDf = spark2 .readStream .format("csv") .schema(csvSchema) .option("header", "true") .option("delimiter", ",") .option("maxFilesPerTrigger", 1) .csv(path)
额外注意
你提供的CSV样本存在表头和数据字段数不匹配的问题:表头共9个字段,每行数据按逗号分隔后有10个值,Schema推断完成后还需要修正CSV格式避免解析异常。
内容的提问来源于stack exchange,提问作者Eljah
相关产品推荐
相关产品推荐

