AWS Glue中DataFrame构造器警告的原因及消除方法咨询
Glue读取S3数据时出现DataFrame构造器内部警告的原因及解决方法
原代码与警告信息
读取S3数据的代码
## 从Amazon S3读取数据 def readFromRawBucket(bucket_name, bucket_prefix, schema_name, table_name): return glue_context.create_dynamic_frame_from_options( connection_type = 's3', connection_options = { 'paths': [f's3://{bucket_name}/{bucket_prefix}/{schema_name}/{table_name}'], 'groupFiles': 'none', 'recurse': True }, format = 'parquet', transformation_ctx = table_name ).toDF()
触发的警告
/home/glue_user/spark/python/pyspark/sql/dataframe.py:127: UserWarning: DataFrame constructor is internal. Do not directly use it. warnings.warn("DataFrame constructor is internal. Do not directly use it.")
警告产生的原因
Glue的DynamicFrame.toDF()方法底层直接调用了Spark内部的DataFrame构造函数,而该构造函数属于Spark内部实现逻辑,并未对外公开推荐使用。Spark官方设计要求用户通过公开的标准API(如spark.read系列方法)创建DataFrame,因此直接调用内部构造器会触发这个警告。
消除警告的修改方案
方案一:通过Spark Session转换DynamicFrame
保留DynamicFrame的处理逻辑,改用Spark Session的公开API将其转换为DataFrame:
## 从Amazon S3读取数据 def readFromRawBucket(bucket_name, bucket_prefix, schema_name, table_name): dynamic_frame = glue_context.create_dynamic_frame_from_options( connection_type = 's3', connection_options = { 'paths': [f's3://{bucket_name}/{bucket_prefix}/{schema_name}/{table_name}'], 'groupFiles': 'none', 'recurse': True }, format = 'parquet', transformation_ctx = table_name ) # 基于DynamicFrame的RDD和Schema,用Spark Session创建DataFrame return glue_context.spark_session.createDataFrame(dynamic_frame.rdd, dynamic_frame.schema())
方案二:直接使用Spark原生读取API
如果不需要DynamicFrame的特定功能,可以绕开它,直接用Spark Session读取S3上的Parquet文件:
## 从Amazon S3读取数据 def readFromRawBucket(bucket_name, bucket_prefix, schema_name, table_name): file_path = f's3://{bucket_name}/{bucket_prefix}/{schema_name}/{table_name}' return glue_context.spark_session.read.parquet(file_path)
方案三:使用Glue Source API
用Glue的getSource方法获取数据源,再转换为DataFrame:
## 从Amazon S3读取数据 def readFromRawBucket(bucket_name, bucket_prefix, schema_name, table_name): source = glue_context.getSource( connection_type='s3', paths=[f's3://{bucket_name}/{bucket_prefix}/{schema_name}/{table_name}'], format='parquet' ) dynamic_frame = source.getFrame() return glue_context.spark_session.createDataFrame(dynamic_frame.rdd, dynamic_frame.schema())
内容的提问来源于stack exchange,提问作者DarkLeafyGreen
相关产品推荐
相关产品推荐

