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

PySpark转换Pandas浮点Series为Arrow字符串数组时出错,Delta表数据写入丢失

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.14 08:53:00