PySpark动态数据框与Pandas数据框的选择/过滤操作差异及AWS Glue场景下的实现咨询
PySpark动态数据框与Pandas数据框的选择/过滤操作差异及AWS Glue场景下的实现咨询
嗨,首先明确告诉你:完全不需要把PySpark/Glue动态数据框转换成Pandas数据框来做基础的选择、过滤操作。反而在AWS Glue这种分布式大数据场景下,转成Pandas会浪费Spark的分布式处理能力,甚至如果数据量超过单节点内存,直接会导致OOM(内存溢出)问题,得不偿失。
接下来给你详细讲下两者的差异,以及在Glue里怎么实现你要的操作:
一、Pandas和PySpark/Glue动态数据框的核心差异
架构与处理能力
- Pandas是单机内存型框架,所有数据都加载在单台机器的内存中,适合小数据集处理;
- PySpark/Glue动态数据框是分布式集群型框架,数据分散存储在集群节点上,能处理TB级别的大数据,这也是Glue作为ETL工具的核心优势。
执行模式
- Pandas是立即执行(Eager Execution):写一行代码就立刻执行计算,结果马上返回;
- PySpark是延迟执行(Lazy Execution):你写的选择、过滤代码只是构建一个计算逻辑的DAG(有向无环图),只有调用
show()、count()、write()这类触发action的方法时,才会真正开始计算。
语法风格
- Pandas偏向于用方括号
[]、点.来直接操作列,语法更简洁; - PySpark/Glue则更偏向于调用API方法(比如
select()、filter())或者用SQL语法,同时支持列表达式的方式。
- Pandas偏向于用方括号
二、AWS Glue场景下的对应实现代码
你可以选择两种方式:直接用Glue的DynamicFrame,或者转换成Spark DataFrame来操作(后者API更丰富灵活,推荐)。
方式1:转换成Spark DataFrame操作(推荐)
先把Glue的DynamicFrame转成Spark DataFrame,然后用Spark的API实现你的需求:
from pyspark.context import SparkContext from awsglue.context import GlueContext from pyspark.sql.functions import concat_ws ctx = SparkContext.getOrCreate() glue_ctx = GlueContext(ctx) # 读取CSV为Glue DynamicFrame,记得开启表头识别 dynamic_frm = glue_ctx.create_dynamic_frame_from_options( connection_type='s3', connection_options={'paths': ['s3://.../abc.csv']}, format='csv', format_options={"withHeader": True} ) # 转成Spark DataFrame spark_df = dynamic_frm.toDF() # 1. 合并first_name和last_name成combined列,等价于Pandas的拼接逻辑 spark_df = spark_df.withColumn('combined', concat_ws(' ', spark_df.first_name, spark_df.last_name)) # 2. 过滤martial_status为married的数据,等价于Pandas的过滤逻辑 filter_result = spark_df.filter(spark_df.martial_status == 'married') # 触发计算查看结果(延迟执行模式下的action操作) filter_result.show() # 如果需要转回DynamicFrame用于Glue后续ETL操作(比如写入Catalog或S3) filtered_dynamic_frm = glue_ctx.create_dynamic_frame_from_rdd(filter_result.rdd, glue_ctx.spark_session)
方式2:直接用Glue DynamicFrame操作
如果不想转Spark DF,也可以用Glue DynamicFrame的原生API:
from pyspark.context import SparkContext from awsglue.context import GlueContext ctx = SparkContext.getOrCreate() glue_ctx = GlueContext(ctx) dynamic_frm = glue_ctx.create_dynamic_frame_from_options( connection_type='s3', connection_options={'paths': ['s3://.../abc.csv']}, format='csv', format_options={"withHeader": True} ) # 1. 合并列:用map方法处理每行数据 def combine_names(rec): rec['combined'] = f"{rec['first_name']} {rec['last_name']}" return rec combined_dynamic_frm = dynamic_frm.map(f=combine_names) # 2. 过滤数据:用filter方法 filtered_dynamic_frm = combined_dynamic_frm.filter(f=lambda rec: rec['martial_status'] == 'married') # 查看结果 filtered_dynamic_frm.show()
三、额外提醒
- 在Glue中读取CSV时,一定要加上
format_options={"withHeader": True},否则列名会被默认命名为_col0、_col1,后续操作会很麻烦; - 如果遇到数据类型不匹配的问题,可以用DynamicFrame的
resolveChoice方法来处理,比如dynamic_frm.resolveChoice(specs=[('martial_status', 'cast:string')]); - 尽量避免在Glue中使用Pandas,除非是处理非常小的数据集(比如配置文件、汇总结果),否则会完全失去分布式处理的优势。
备注:内容来源于stack exchange,提问作者nohardfeelings
相关产品推荐
相关产品推荐

