You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

调用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()

关键修改点说明

  1. SparkSession初始化:移除原代码中sc = SparkContext()的显式创建,改用SparkSession.builder.getOrCreate()获取会话,再基于会话的sparkContext创建GlueContext,完全遵循弃用提示要求,解决警告问题。
  2. 变量错误修正:原代码中dbtable引用未定义的db_table,改为使用函数传入的table参数,避免变量未定义错误。
  3. 方法名修正:将to_DF()改为Spark/Glue正确的方法名toDF(),此拼写错误可能是触发awaitResult异常的原因之一。

内容的提问来源于stack exchange,提问作者Bharath M

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.30 20:24:32