AWS Lambda操作Google CloudSQL Postgres:SELECT正常但INSERT无效果
问题:CloudSQL INSERT操作无报错但不生效,Query Insights显示已执行但目标表为空
连接Google CloudSQL的服务账号认证、数据库用户名及密码均正常,SELECT语句可正常执行,但INSERT操作无报错、无异常却无法生效。后续通过CloudSQL的Query insights确认INSERT查询已执行,但目标表zoom_account仍为空。
相关Python代码
def execute_statement(db: sqlalchemy.engine.base.Engine, rows) -> None: start_time = time.time() invoice_month = str(now.month) invoice_period = "%s%s" % (now.year, invoice_month.zfill(2)) sql_query = text(""" INSERT INTO zoom_account (id, month, account_number, service, cost, updated) VALUES (:id, :month, :account_number, :service, :cost, :updated) ON CONFLICT (id, month) DO UPDATE SET account_number = EXCLUDED.account_number, service = EXCLUDED.service, cost = EXCLUDED.cost, updated = EXCLUDED.updated RETURNING id """) data = ( {"id": str(hash(row.account_num + ", " + row.service_name) % 10000), "month": invoice_period, "account_number": row.account_num, "service": row.service_name, "cost": Decimal(row.estimated_cost), "updated": datetime.datetime.now()} for row in rows ) with db.connect() as conn: try: for line in data: conn.execute(sql_query, line) except SQLAlchemyError as e: logger.error(f"Error inserting or updating data into database: {str(e)}") # query = text("SELECT * FROM zoom_account") # with db.connect() as conxn: # res = conxn.execute(query).fetchall() # logger.info(f'Result: {res}') # This select works without issue duration = time.time() - start_time logger.info(f"{len(rows)} Insert or update took {duration:.2f} seconds.")
Lambda运行日志
2023-03-06T11:35:05.206+11:00 2023-03-06 00:35:05,206 - extract - INFO - Initialising database connection ... 2023-03-06T11:35:39.251+11:00 2023-03-06 00:35:39,250 - extract - INFO - 106 Insert or update took 31.11 seconds.
解决方案
- 事务未提交是核心问题:SQLAlchemy通过
db.connect()获取连接时,默认会开启一个事务,执行完execute后如果没有手动提交,事务会在连接关闭时自动回滚,导致数据不会写入数据库。
修正方法:在循环执行完插入/更新后添加提交操作,或者使用begin()上下文管理器自动处理事务:
方法1:手动提交
方法2:使用with db.connect() as conn: try: for line in data: conn.execute(sql_query, line) conn.commit() # 新增事务提交语句 except SQLAlchemyError as e: conn.rollback() # 异常时回滚事务 logger.error(f"Error inserting or updating data into database: {str(e)}")begin()自动管理事务(更简洁)with db.begin() as conn: # begin()会在块结束时自动提交,异常时自动回滚 try: for line in data: conn.execute(sql_query, line) except SQLAlchemyError as e: logger.error(f"Error inserting or updating data into database: {str(e)}") - 验证执行结果:可以捕获
execute的返回值,打印RETURNING的结果,确认操作是否真的生效:result = conn.execute(sql_query, line) returned_id = result.fetchone() logger.info(f"Processed row, returned id: {returned_id}") - 排查冲突逻辑:虽然表为空时
ON CONFLICT不会触发,但可以检查生成的id和month组合是否存在重复(比如哈希碰撞),不过这不是当前表为空的主要原因。
内容的提问来源于stack exchange,提问作者PainPoints
相关产品推荐
相关产品推荐

