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

如何在AWS Glue作业中用PySpark SQL更新PostgreSQL表列

在AWS Glue中基于Square客户数据实现square_dashboard表的UPSERT操作

我通过AWS Lambda将Square数据同步到PostgreSQL,再用AWS Glue做ETL。现有两个Glue作业分别处理订单和客户数据,现在需要基于square_customer表的聚合结果,对square_dashboard表做更新现有行+插入新行的操作。我已经有能在PostgreSQL中正常运行的SQL语句,但用Glue的spark.sql()执行时遇到连接和语法兼容问题,求可行的实现方案。

原PostgreSQL SQL语句:

-- 更新现有行
UPDATE square_dashboard sd
SET new_customer = subquery.new_customers
FROM (
    SELECT user_id, DATE(created_at) AS customer_date,
           COUNT(*) AS new_customers
    FROM square_customer
    GROUP BY user_id, DATE(created_at)
) AS subquery
WHERE sd.user_id = subquery.user_id
  AND sd.date = subquery.customer_date;

-- 插入新行
INSERT INTO square_dashboard (user_id, date, new_customer, day)
SELECT subquery.user_id, 
       subquery.customer_date, 
       subquery.new_customers,
       TO_CHAR(subquery.customer_date, 'Dy')
FROM (
    SELECT user_id, DATE(created_at) AS customer_date,
           COUNT(*) AS new_customers
    FROM square_customer
    GROUP BY user_id, DATE(created_at)
) AS subquery
WHERE NOT EXISTS (
    SELECT 1 
    FROM square_dashboard sd 
    WHERE sd.user_id = subquery.user_id
      AND sd.date = subquery.customer_date
);

方案一:用JDBC直接执行原生PostgreSQL SQL

spark.sql()默认操作的是Spark元数据中的表(Glue Data Catalog或临时视图),无法直接执行PostgreSQL原生的UPDATE/INSERT语法。可以通过Python的JDBC驱动直接连接PostgreSQL,执行你已有的SQL语句。

实现步骤

  1. 确保Glue作业的执行角色拥有PostgreSQL的VPC访问权限、数据库凭证读取权限。
  2. 为Glue作业添加psycopg2-binary依赖(可通过Glue层或自定义Python库路径实现)。
  3. 编写代码连接数据库并执行SQL。

代码示例

import psycopg2
from awsglue.utils import getResolvedOptions
import sys

# 从Glue作业参数读取数据库配置(推荐用Secrets Manager或Glue连接存储凭证,避免硬编码)
args = getResolvedOptions(sys.argv, ['JOB_NAME', 'DB_HOST', 'DB_PORT', 'DB_NAME', 'DB_USER', 'DB_PASSWORD'])

# 建立PostgreSQL连接
conn = psycopg2.connect(
    host=args['DB_HOST'],
    port=args['DB_PORT'],
    dbname=args['DB_NAME'],
    user=args['DB_USER'],
    password=args['DB_PASSWORD']
)
cur = conn.cursor()

# 执行更新语句
update_sql = """
UPDATE square_dashboard sd
SET new_customer = subquery.new_customers
FROM (
    SELECT user_id, DATE(created_at) AS customer_date,
           COUNT(*) AS new_customers
    FROM square_customer
    GROUP BY user_id, DATE(created_at)
) AS subquery
WHERE sd.user_id = subquery.user_id
  AND sd.date = subquery.customer_date;
"""
cur.execute(update_sql)

# 执行插入语句
insert_sql = """
INSERT INTO square_dashboard (user_id, date, new_customer, day)
SELECT subquery.user_id, 
       subquery.customer_date, 
       subquery.new_customers,
       TO_CHAR(subquery.customer_date, 'Dy')
FROM (
    SELECT user_id, DATE(created_at) AS customer_date,
           COUNT(*) AS new_customers
    FROM square_customer
    GROUP BY user_id, DATE(created_at)
) AS subquery
WHERE NOT EXISTS (
    SELECT 1 
    FROM square_dashboard sd 
    WHERE sd.user_id = subquery.user_id
      AND sd.date = subquery.customer_date
);
"""
cur.execute(insert_sql)

# 提交事务并关闭连接
conn.commit()
cur.close()
conn.close()

方案二:用PySpark DataFrame实现UPSERT(适合大数据场景)

对于大规模数据处理,推荐用Spark分布式能力实现UPSERT逻辑,避免单节点JDBC执行的性能瓶颈。

实现步骤

  1. 从PostgreSQL读取square_customer和square_dashboard为DataFrame。
  2. 聚合计算每日新增客户数。
  3. 拆分更新、插入数据集,分别执行操作;或用PostgreSQL的ON CONFLICT实现原子UPSERT。

代码示例

from awsglue.context import GlueContext
from pyspark.context import SparkContext
from pyspark.sql.functions import count, date_format, col, to_date

# 初始化Glue上下文
sc = SparkContext()
glueContext = GlueContext(sc)
spark = glueContext.spark_session

# PostgreSQL连接配置(推荐用Glue连接管理凭证)
jdbc_url = "jdbc:postgresql://your-db-host:5432/your-db-name"
connection_properties = {
    "user": "your-db-user",
    "password": "your-db-password",
    "driver": "org.postgresql.Driver"
}

# 读取客户数据并聚合
customer_df = spark.read.jdbc(url=jdbc_url, table="square_customer", properties=connection_properties)
agg_customer_df = customer_df.groupBy("user_id", to_date(col("created_at")).alias("customer_date")) \
    .agg(count("*").alias("new_customers")) \
    .withColumn("day", date_format(col("customer_date"), "EEE"))  # 对应PostgreSQL的TO_CHAR('Dy')

# 方式1:拆分更新和插入
# 读取仪表板数据
dashboard_df = spark.read.jdbc(url=jdbc_url, table="square_dashboard", properties=connection_properties)

# 生成更新数据集(匹配已有记录)
update_df = agg_customer_df.join(dashboard_df, 
                                 (agg_customer_df.user_id == dashboard_df.user_id) & 
                                 (agg_customer_df.customer_date == dashboard_df.date),
                                 "inner") \
    .select(dashboard_df["id"],  # 假设id是仪表板表主键
            agg_customer_df["new_customers"].alias("new_customer"))

# 执行更新
update_df.createOrReplaceTempView("temp_updates")
spark.sql("""
    UPDATE square_dashboard 
    SET new_customer = t.new_customer
    FROM temp_updates t
    WHERE square_dashboard.id = t.id
""")

# 生成插入数据集(无匹配记录)
insert_df = agg_customer_df.join(dashboard_df, 
                                 (agg_customer_df.user_id == dashboard_df.user_id) & 
                                 (agg_customer_df.customer_date == dashboard_df.date),
                                 "left_anti") \
    .select(col("user_id"),
            col("customer_date").alias("date"),
            col("new_customers").alias("new_customer"),
            col("day"))

# 执行插入
insert_df.write.jdbc(url=jdbc_url, 
                     table="square_dashboard", 
                     mode="append", 
                     properties=connection_properties)

# 方式2:用PostgreSQL ON CONFLICT实现原子UPSERT(更高效)
agg_customer_df.select(
    col("user_id"),
    col("customer_date").alias("date"),
    col("new_customers").alias("new_customer"),
    col("day")
).write.jdbc(
    url=jdbc_url,
    table="square_dashboard",
    mode="append",
    properties={
        **connection_properties,
        "statement": """
            INSERT INTO square_dashboard (user_id, date, new_customer, day)
            VALUES (?, ?, ?, ?)
            ON CONFLICT (user_id, date) DO UPDATE SET new_customer = EXCLUDED.new_customer
        """
    }
)

内容的提问来源于stack exchange,提问作者Muhammad Ahmed

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 03:43:10