调用o98.getDynamicFrame报错,按SparkSession重构代码未解决
问题解决:Glue作业中SparkSession弃用警告与代码重构
问题现象
- 调用
o98.getDynamicFrame时触发awaitResult异常 - 尝试用
SparkSession.builder.getOrCreate()重构代码后,仍出现弃用警告:
FutureWarning: Deprecated in 3.0.0. Use SparkSession.builder.getOrCreate() instead.
修正后的完整代码
import logging import sys from awsglue.job import Job from awsglue.utils import getResolvedOptions from awsglue.context import GlueContext from pyspark.sql import SparkSession from pyspark.sql.functions import when, lit # 日志配置 logger = logging.getLogger(__name__) handler = logging.StreamHandler(sys.stdout) handler.setLevel(logging.INFO) formatter = logging.Formatter('%(asctime)s - %(name)s - %(levelname)s - %(message)s') handler.setFormatter(formatter) logger.addHandler(handler) # 获取作业参数 args = getResolvedOptions(sys.argv, ['JOB_NAME', 'TempDir']) USER_NAME = 'user' PASSWORD = 'password' # 正确初始化SparkSession与GlueContext spark = SparkSession.builder.getOrCreate() glueContext = GlueContext(spark.sparkContext) # 初始化Glue作业 job = Job(glueContext) job.init(args['JOB_NAME'], args) def get_df(table: str): conn_options = { "url": "redshift-url", "user": USER_NAME, "password": PASSWORD, "dbtable": table, # 修正:使用函数参数table替代未定义的db_table "redshiftTmpDir": args["TempDir"], "aws_iam_role": "role" } # 修正:to_DF()改为正确的toDF() return glueContext.create_dynamic_frame_from_options("redshift", conn_options).toDF()
关键修改点说明
- SparkSession初始化:移除原代码中
sc = SparkContext()的显式创建,改用SparkSession.builder.getOrCreate()获取会话,再基于会话的sparkContext创建GlueContext,完全遵循弃用提示要求,解决警告问题。 - 变量错误修正:原代码中
dbtable引用未定义的db_table,改为使用函数传入的table参数,避免变量未定义错误。 - 方法名修正:将
to_DF()改为Spark/Glue正确的方法名toDF(),此拼写错误可能是触发awaitResult异常的原因之一。
内容的提问来源于stack exchange,提问作者Bharath M
相关产品推荐
相关产品推荐

