PySpark转换Pandas浮点Series为Arrow字符串数组时出错,Delta表数据写入丢失
我来帮你一步步排查和解决这个问题,先从最明显的错误入手,再处理类型转换和Delta Schema的核心矛盾:
1. 先修复函数调用的笔误(数据丢失的关键原因)
看你的process_and_write_page函数里,线程池调用的是fetch_location_details,但你实际定义的获取详情的函数是fetch_details——这个笔误会导致每个线程都抛出NameError,根本拿不到数据!这也是为什么最终只有11000条记录的核心原因之一,先把这个改过来:
# 原错误代码 # future_to_id = {executor.submit(fetch_location_details, id): id for id in ids} # 修正为 future_to_id = {executor.submit(fetch_details, id): id for id in ids}
2. 解决Pandas转Spark的类型不匹配错误
你遇到的Exception thrown when converting pandas.Series (float64) with name 'latitude' to Arrow Array (string)错误,本质是类型校验不兼容:
- Pandas里
latitude列是float64类型(API返回的数值被转成Python float,Pandas自动识别为浮点型) - 但你定义的Spark Schema里
latitude是StringType - 开启Arrow优化后,PySpark会严格校验类型,不允许直接把浮点类型强制转成字符串类型
修复方案:统一类型,提前处理异常值
推荐方案:将Spark Schema的纬度/经度改为FloatType
首先修改basic_schema的类型定义:
basic_schema = StructType([ StructField("id", StringType(), True), StructField("type", StringType(), True), StructField("name", StringType(), True), StructField("latitude", FloatType(), True), StructField("longitude", FloatType(), True), StructField("alsoKnownAs", StringType(), True) ])
然后在fetch_details函数里,确保返回的纬度/经度是数值类型,同时处理API可能返回的空值或非数值字符串:
# 原代码 # "latitude": data.get("latitude"), # "longitude": data.get("longitude"), # 修正为 "latitude": float(data.get("latitude")) if data.get("latitude") not in (None, "") else None, "longitude": float(data.get("longitude")) if data.get("longitude") not in (None, "") else None,
最后在转换为Spark DataFrame前,在Pandas层面再做一次数值校验,把无法转换的值转为NaN(Spark会识别为null):
if not basic_df.empty: # 处理Pandas中的非数值数据,errors='coerce'会把无效值转成NaN basic_df['latitude'] = pd.to_numeric(basic_df['latitude'], errors='coerce') basic_df['longitude'] = pd.to_numeric(basic_df['longitude'], errors='coerce') try: basic_spark_df = spark.createDataFrame(basic_df, schema=basic_schema) basic_spark_df.write.mode("append").option("mergeSchema", "true").format("delta").saveAsTable("api_test") except Exception as e: print(f"处理页面数据失败,ID列表:{ids},错误信息:{e}") # 可选:把错误数据保存到本地文件用于排查 basic_df.to_csv(f"error_page_{p}.csv", index=False)
3. 解决Delta Schema合并失败的问题
你改成FloatType时遇到的AnalysisException: [DELTA_FAILED_TO_MERGE_FIELDS],是因为你的Delta表已经存在latitude为StringType的历史数据——Delta不允许直接把字符串类型改为浮点类型(类型不兼容)。
修复方案二选一:
方案A:重新初始化Delta表(适合测试阶段,数据可重新获取)
先删除现有表,再重新运行代码:
# 运行一次删除表(之后记得注释掉) spark.sql("DROP TABLE IF EXISTS api_test")
方案B:迁移现有数据到新表(保留历史数据)
如果不想丢失已有的11000条数据,可以创建新表转换类型后替换旧表:
# 读取旧表,转换latitude/longitude类型 old_df = spark.read.table("api_test") new_df = old_df.withColumn("latitude", old_df["latitude"].cast(FloatType())) \ .withColumn("longitude", old_df["longitude"].cast(FloatType())) # 写入新表 new_df.write.mode("overwrite").format("delta").saveAsTable("api_test_new") # 替换旧表 spark.sql("DROP TABLE IF EXISTS api_test") spark.sql("ALTER TABLE api_test_new RENAME TO api_test")
之后再运行数据获取代码,就不会出现Schema合并错误了。
4. 优化错误处理,避免页面级数据丢失
原来的代码中,如果createDataFrame出错,整个页面的数据都不会写入。我们已经在上面的代码中添加了try-except块,这样即使单个页面有错误,也只会跳过该页面(或保存错误数据排查),不会影响其他页面的写入。
备注:内容来源于stack exchange,提问作者CrowsNose

