如何在Polars中并行化返回字符串类型的H3 Polyfill Pandas UDF
解决方案:Polars中实现H3 Polyfill并行化处理
针对你遇到的Polars与地理数据、H3交互的问题,以下是分步解决思路及完整实现代码:
核心问题分析
- Geo类型兼容性:Polars目前原生不支持GeoArrow/GeoParquet类型,因此直接转换GeoDataFrame会报错,需从原始WKT字符串入手处理。
- Polars Expr与GeoPandas不兼容:GeoPandas的方法无法直接接收Polars表达式,需在Polars的UDF(用户定义函数)或批量处理接口中封装逻辑。
- Object类型无法Explode:
apply返回的Python对象列表会被识别为Object类型,需明确转为Polars的List<Utf8>类型才能执行explode操作。
完整实现代码
1. 导入依赖库
import polars as pl from shapely import wkt import h3 import pandas as pd
2. 定义并行化UDF
使用Polars的UDF装饰器,明确输入输出类型,确保并行处理及类型兼容性:
@pl.api.udf(input_types=[pl.Utf8], return_type=pl.List(pl.Utf8)) def h3_polyfill_from_wkt(wkt_str: str, res: int = 10) -> list[str]: # 解析WKT为几何对象 geom = wkt.loads(wkt_str) # 转换为GeoJSON格式供H3处理 geojson = wkt.geometry.mapping(geom) # 执行Polyfill并转为字符串列表 return list(h3.polyfill(geojson, res=res, geo_json_conformant=True))
3. 处理数据集并生成结果
# 原始数据集(你的输入数据) pandas_df = pd.DataFrame({ 'quadkey': {0: '0022133222330023', 1: '0022133222330031', 2: '0022133222330100'}, 'tile': { 0: 'POLYGON((-160.043334960938 70.6363054807905, -160.037841796875 70.6363054807905, -160.037841796875 70.6344840663086, -160.043334960938 70.6344840663086, -160.043334960938 70.6363054807905))', 1: 'POLYGON((-160.032348632812 70.6381267305321, -160.02685546875 70.6381267305321, -160.02685546875 70.6363054807905, -160.032348632812 70.6363054807905, -160.032348632812 70.6381267305321))', 2: 'POLYGON((-160.02685546875 70.6417687358462, -160.021362304688 70.6417687358462, -160.021362304688 70.6399478155463, -160.02685546875 70.6399478155463, -160.02685546875 70.6417687358462))' }, 'avg_d_kbps': {0: 15600, 1: 6790, 2: 9619}, 'avg_u_kbps': {0: 14609, 1: 22363, 2: 15757}, 'avg_lat_ms': {0: 168, 1: 68, 2: 92}, 'tests': {0: 2, 1: 1, 2: 6}, 'devices': {0: 1, 1: 1, 2: 1} }) # 转换为Polars DataFrame polars_df = pl.from_pandas(pandas_df) # 添加H3单元格列,保留所有原始列 result_df = polars_df.with_columns( h3_polyfill_from_wkt(pl.col('tile'), res=10).alias('h3_cells') ) # 展开H3单元格列表,完成类似h3pandas的resample效果 exploded_df = result_df.explode('h3_cells').collect() # 查看结果 print(exploded_df)
性能优化可选方案
如果处理超大规模数据集,可改用map_batches批量处理,减少逐元素调用的开销:
def polyfill_batch(wkt_series: pl.Series) -> pl.Series: # 批量解析WKT geometries = wkt.loads(wkt_series.to_pandas()) # 批量生成H3列表 h3_lists = [list(h3.polyfill(wkt.geometry.mapping(g), res=10, geo_json_conformant=True)) for g in geometries] # 返回Polars List类型序列 return pl.Series(h3_lists, dtype=pl.List(pl.Utf8)) # 应用批量处理 result_df = polars_df.with_columns( pl.col('tile').map_batches(polyfill_batch).alias('h3_cells') )
关键解决点说明
- 绕过GeoPandas的几何类型,直接从WKT字符串处理,避免Polars的类型兼容性问题。
- 明确指定UDF的返回类型为
pl.List(pl.Utf8),确保Polars识别为列表类型,支持explode操作。 - Polars的UDF和
map_batches默认利用多核并行处理,性能远高于Pandas的单线程apply,适合大规模数据集。
内容的提问来源于stack exchange,提问作者Georg Heiler
相关产品推荐
相关产品推荐

