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

Spark 3.3.0结构化流写入时部分CSV行丢失问题求助

Spark结构化流写入时部分记录丢失问题(疑似字符过长导致)

环境信息

  • Spark版本:Apache Spark 3.3.0

问题描述

使用Spark结构化流读取并处理逗号分隔的CSV文件,读取后的DataFrame数据正常,但写入目标表时部分记录自动丢失,推测是某列值包含大量字符导致。

丢失记录示例

以下是一条丢失的逗号分隔行:

a500o0000008bugAAA,FALSE,KMI000004704,Key Medical Insight,0050o00000WuoSBAAZ,2020-04-02T10:17:02.000Z,0019000000R3GVDAA3,4/2/2020,"<XXXX XXXXX=XXXXX-XXXXX: XXX-XXXX;>XXXXXXXXXX XXXX XXXXXXX XXXX XX XXXXXXXXXXX XXXXXXXX XXXXXX XXXXX-XX. XXXXXX XXX XXX XXXXXXXXX XXX XXXXX XXXXXXXXXX XXXX XXXXXXX XXXXXXXXXXX XXX XXXX XXXXXXXX XX XXXXXXXXXX XXXXX XXXXXXXXXXXX XXXXX XX XXXXX XXXXX. XXXXXXXX XXXX XXXX XX XXXXX XXXXX XXXXXXXXXX XXXXXX XXXXX-XX XXXXX XXXXXXXX XXXXXX XXXXXXXXXX XXX XXXXXXX XXXX XXXXXXX XXXXXX XX XXXX XXXXXX XXXXXXXXXX XXX XXXXXXX XXXXX XXX. </XXXX><XXX><XXXX XXXXX=XXXXX-XXXXX: XXX-XXXX;><XX></XXXX></XXX><XXX><XXXX XXXXX=XXXXX-XXXXX: XXX-XXXX;>XXXXXXXXX XXXX XXXXXXX&#XX;X XX.X XXXXX XXXXXXXXX (XXX-XXXXXXXX) XXXXX XXX XXXX XXXXXXXX XXXX X XXXXXXXX - XXXX XXXXXXXX XX XXXXXX XXX (XXXXXX XXXX XXXXX XXX XXXXX) XX XXXX-XX-XXXX XXXXX XXX XXXXXXX XXXXXXX XXXXXXXXX XXX XXXXXXXXXX (XXXX XXX XXXXXXXX XXXXX). 

XXXXXXX XXXXXXXXX XXXX XXXXXXXX XX XXXXXX XXX XXXXXX XXXX XXXX XX XXXXX XXXXXXX X.X. &XXXX;XXXXXXXX XXXXX XXXXXXXX XXXXX XXXXXXXXXX XXXX XXXX XXXX XXXXXXX XXXXXXXXX&XXXX;. </XXXX><XXX><XXXX XXXXX=XXXXX-XXXXX: XXX-XXXX;><XX></XXXX></XXX><XXX><XXXX XXXXX=XXXXX-XXXXX: XXX-XXXX;>XXXXX XXXXXXXXXX: XXXXXXX XXXXXXX XX XX XXXXX XXXXXXXXXXX XXXXX XXXXXXXXXX XX XXXX XXXXXXXXXX XX XXXXXXXXXX XXXXXXXX XXXXX XXX XXXXXXX XX XXXXXXXXX XXXXXXXXXX XXXXXXXXXXXX XX XXXX XXXXXXXXXX XXXXX XX XXXXXXXX XXXXXXXX XXXX XXX XX X XXXXX XXXX XXXX XXX XXXXXX XXXXXXXXXX XX XXXXXXXXXX XXXXX XXXXXXXXX. XX XXXXXXXX XXXX XXXXXXXX XXXXXXXXXXXX XXXXXXXXX XXX XXXXX XXXXXXXXX XXXXXX XXXXX-XX XXXXXX XXXXXXX XX XX XXXXX XXXX XXXXXXXXX XXXXXXXX XXX XXXXXXXXXX/XXXXXXX - XXXXXX XX XXX.</XXXX></XXX></XXX>",a040o00002QZ26mAAD,Submitted_vod,Local Data/ Fact/Observation,Key Opinion Leader,Hematology

注:包含"xxxxx"的字段是一个包含大量空格和特殊字符的单个字段。

读取CSV文件的代码

def read_stream(container_read_path, file_format, delimeter, spark, header):   
 spark.conf.set("spark.sql.streaming.schemaInference", True)
 source_data = (
    spark.readStream.format(file_format)
    .option("header", header)
    .option("sep", delimeter)
    .option("escape", "\"")
    .option("multiline", True)
    .option("recursiveFileLookup", "true")
    .load(f"{container_read_path}")
 )    
 return source_data

写入结构化流数据的代码

def write_stream(dataframe, database_name, table_name, checkpoint_path, partition_cols, 
header, file_format='parquet'):

 (dataframe
  .writeStream
  .format(file_format)
  .trigger(once=True)
  .option("checkpointLocation", f'{checkpoint_path}')
  .foreachBatch(lambda df, epochId: write_raw_file(df, epochId, database_name, 
   table_name, partition_cols, header, file_format))
  .start()
  )

def write_raw_file(df, epochId, database_name, table_name, partition_cols, header, 
file_format):
 file_format = 'csv' if file_format == 'text' else file_format
 header = "true" if file_format == 'csv' else "false"

 (df.write
  .mode("append")
  .option("header", header)
  .partitionBy(partition_cols)
  .format(file_format)
  .saveAsTable(f"{database_name}.{table_name}")
  )

内容的提问来源于stack exchange,提问作者Shahrukh lodhi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 09:20:29