如何在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语句。
实现步骤
- 确保Glue作业的执行角色拥有PostgreSQL的VPC访问权限、数据库凭证读取权限。
- 为Glue作业添加
psycopg2-binary依赖(可通过Glue层或自定义Python库路径实现)。 - 编写代码连接数据库并执行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执行的性能瓶颈。
实现步骤
- 从PostgreSQL读取
square_customer和square_dashboard为DataFrame。 - 聚合计算每日新增客户数。
- 拆分更新、插入数据集,分别执行操作;或用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
相关产品推荐
相关产品推荐

