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

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:手动提交
    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)}")
    
    方法2:使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 03:35:35