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

能否将大型Spark DataFrame拆分至Hive表后循环处理生成Folium地图?

关于Spark DataFrame生成Folium地图的解决方案

首先明确说:把大数据拆分到Hive表再循环处理是可行的,但这并不是最高效的解决方案——咱们来拆解下问题,再看看更优的处理方式:

一、拆分到Hive表的可行性分析

这种思路本质是将分布式存储的Spark数据,拆分成Hive的分区(或分表),然后逐批读取小批量数据转成Pandas DataFrame,再用你熟悉的循环生成Folium标记。它能解决单批次加载超大数据到内存的问题,但缺点也很明显:多了Hive读写的额外开销,流程相对繁琐。

如果一定要这么做,大致代码逻辑是这样的:

# 1. 将Spark DataFrame写入Hive表(可以按字段分区优化读取)
spark_df.write.partitionBy("pickup_date").saveAsTable("hive_pickup_locations")

# 2. 初始化Spark会话并获取所有分区
from pyspark.sql import SparkSession
spark = SparkSession.builder.enableHiveSupport().getOrCreate()
partitions = spark.sql("SHOW PARTITIONS hive_pickup_locations").collect()

# 3. 逐分区读取并生成Folium标记
marker_cluster = folium.MarkerCluster().add_to(your_map)
for partition in partitions:
    # 读取单个分区的数据
    part_data = spark.sql(f"SELECT Pickup_latitude, Pickup_longitude FROM hive_pickup_locations WHERE {partition[0]}")
    # 转成Pandas小批量处理
    pd_part = part_data.toPandas()
    for _, row in pd_part.iterrows():
        folium.CircleMarker(
            location=(row["Pickup_latitude"], row["Pickup_longitude"]),
            radius=20,
            color="#0A8A9F",
            fill=True
        ).add_to(marker_cluster)

二、更高效的替代方案

其实不用绕Hive,直接在Spark层面处理后转成Pandas小批量数据即可,核心思路是先过滤/采样,再分批次转换:

1. 先做数据预处理(必做!)

超大数据集全量渲染到Folium地图上会直接导致浏览器崩溃,所以第一步必须过滤无效数据+采样:

# 过滤掉不在合理范围的经纬度,同时采样10%的数据(比例按需调整)
processed_spark_df = spark_df.filter(
    (spark_df["Pickup_latitude"].between(-90, 90)) & 
    (spark_df["Pickup_longitude"].between(-180, 180))
).sample(fraction=0.1, seed=42)

2. 直接转Pandas(如果采样后数据量不大)

如果采样后的数据能塞进内存,直接转成Pandas再循环:

pd_df = processed_spark_df.toPandas()
marker_cluster = folium.MarkerCluster().add_to(your_map)
for _, row in pd_df.iterrows():
    folium.CircleMarker(
        location=(row["Pickup_latitude"], row["Pickup_longitude"]),
        radius=20,
        color="#0A8A9F",
        fill=True
    ).add_to(marker_cluster)

3. 分批次转换(采样后仍过大)

如果采样后数据还是超出内存,可以用randomSplit拆成多个小的Spark DataFrame,逐个处理:

# 拆分成3个批次(数量按需调整)
split_dfs = processed_spark_df.randomSplit([0.33, 0.33, 0.34], seed=42)

marker_cluster = folium.MarkerCluster().add_to(your_map)
for df in split_dfs:
    pd_batch = df.toPandas()
    for _, row in pd_batch.iterrows():
        folium.CircleMarker(
            location=(row["Pickup_latitude"], row["Pickup_longitude"]),
            radius=20,
            color="#0A8A9F",
            fill=True
        ).add_to(marker_cluster)

三、额外提示

如果你的数据量实在大到采样后都没法处理,还可以试试Spark的foreachPartition方法——在每个分布式分区内直接生成Folium标记,最后收集到Driver端合并。不过这种方式需要注意Folium依赖在Executor节点的安装,以及序列化问题,复杂度稍高,一般采样方案足够应对大多数场景。

内容的提问来源于stack exchange,提问作者A.HADDAD

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 08:15:53