Databricks PySpark readStream读ORC时误读非ORC文件报错问题
异常根因
这个报错和ORC文件本身损坏、租约锁阻塞ORC读取没有关系,是Structured Streaming文件源扫描逻辑和ADLS Gen2租约特性不匹配导致的,具体触发逻辑:
- 首先
readStream的文件数据源默认不会按文件后缀做过滤:只要是传入根路径下(含所有loadDate分区子目录)的文件,只要被判定为"已写入完成",就会被加入待处理列表,按你指定的ORC格式解析,不会因为后缀是.json就自动跳过。之前同逻辑作业运行正常,本质是之前作业触发时,路径下的JSON文件还处于写入排他锁状态,被流的未完成文件过滤逻辑跳过了。 - 你提到的JSON文件被其他进程持有租约锁,恰恰是这次异常的触发点:Structured Streaming默认跳过写入中文件的逻辑,依赖文件系统访问时抛出的"文件被占用/不可读"IO异常做判断,但ADLS Gen2的租约分排他写租约、共享读租约两种,如果其他进程持有的是共享租约,DBFS的ADLS连接器可以正常读取文件元数据、文件头内容,不会抛出IO异常,流引擎就会判定这个JSON文件是已写完的有效待处理文件,直接按ORC格式解析,读到JSON格式的文件头就会抛出
org.apache.orc.FileFormatException: Malformed ORC file错误。 - 你配置的
trigger(once=True)模式会放大这个问题:一次性触发模式会在单批次内扫描所有checkpoint未记录的新增文件,不会像常驻流作业那样把识别异常的文件留到后续批次重试,只要待处理列表混进非ORC文件就直接失败。
修复方法
不需要调整现有存储目录结构,在readStream配置里加参数即可彻底解决:
- 加后缀过滤规则,只读取ORC格式的文件,从根源避免扫描到JSON等其他格式文件:
stream_df = spark.readStream \ .format("orc") \ .schema(pre_loaded_schema) \ # 你预加载的表结构 .option("pathGlobFilter", "*.orc") \ # 仅匹配后缀为.orc的文件 .load("/mnt/your_adls_mount/root_path") - 可选优化:如果根目录下除了分区目录、JSON文件还有其他临时文件,可以开启分区目录自动发现,避免扫到根目录下的零散非分区文件:
# 在上边的readStream配置里追加 .option("recursiveFileLookup", "true") - 注意不要为了绕过报错开启
spark.sql.files.ignoreCorruptFiles=true,这个参数会静默跳过所有解析失败的文件,可能把真正损坏的ORC文件也跳过,导致数据漏读。
内容的提问来源于stack exchange,提问作者Tim
相关产品推荐
相关产品推荐

