如何在循环中创建不同名称的Spark DataFrame(不使用Pandas)
问题
我有一个从AWS S3存储的JSON文件创建Spark DataFrame的AWS Glue函数(代码如下),需要遍历S3中的表文件夹列表,在循环中为每个表创建对应名称的DataFrame(如orders对应df_orders)。尝试过字典方法但无法获取Spark DataFrame,现询问不使用Pandas的实现方式。
函数代码
def create_glue_df(table): df = glueContext.create_dynamic_frame.from_options( format_options={"jsonPath": "$._airbyte_data", "multiline": True}, connection_type="s3", format="json", connection_options={ "paths": [f"s3://bucketx/raw-data/test/{table}/"], "recurse": True, }, transformation_ctx=f"{table}", ) df = df.toDF() return df
表列表
table_list = ['orders', 'products', 'price', 'inventory', etc.]
当前循环代码
for table in table_list : try: df = create_glue_df(table) print(f'Table {table} created.') except: print(f'Table {table} not found.')
解决方案
用字典存储每个表对应的Spark DataFrame是最适配Spark生态的实现方式,完全无需依赖Pandas。以表名为字典的key,对应DataFrame为value,后续可通过key直接访问目标DataFrame。
修改后的循环代码如下:
# 初始化空字典用于存储所有DataFrame df_dict = {} for table in table_list : try: df = create_glue_df(table) # 将生成的DataFrame存入字典 df_dict[table] = df print(f'Table {table} created and stored in df_dict["{table}"].') except Exception as e: # 捕获具体异常便于排查问题 print(f'Failed to create table {table}: {str(e)}')
访问DataFrame的方式
后续可通过字典键调用对应表的DataFrame,示例:
# 获取orders表的DataFrame df_orders = df_dict['orders'] # 查看数据结构 df_orders.printSchema()
关于之前字典方法失败的可能原因
如果之前尝试字典方式未成功,大概率是未正确将生成的DataFrame赋值到字典中,或是后续访问方式有误。上述代码明确将create_glue_df返回的Spark DataFrame存入字典,只要函数本身能正常返回DataFrame,即可实现正确存储与访问。
内容的提问来源于stack exchange,提问作者daniellidio
相关产品推荐
相关产品推荐

