MySQL更新表时触发1412错误:表定义已变更,请重试事务
问题:MySQL事务中更新关联临时表时触发1412错误
开发个人项目时,流程为:创建并加载临时表,关联数据库现有表,基于临时表结果更新现有表。执行db.execute_insert_update()时抛出异常:
mysql.connector.errors.DatabaseError: 1412 (HY000): Table definition has changed, please retry transaction
已尝试为临时表指定schema,但问题依旧。相关代码如下:
def _create_and_load_table(self, dataframe: DataFrame) -> str: #first_function self._log_entry() temp_table_name = self._generate_table_name() table_schema = "partner_staging" table_with_schema = f"{table_schema}.{temp_table_name}" self.df_mgr.jdbc_write(dataframe=dataframe, dbtable=table_with_schema, check_count_before_write=True) self._log_exit() return temp_table_name def reconcile_partner_payments(self, dataframe: DataFrame) -> NoReturn: #second function self.logger.info("Start reconcile_partner_payments") self.logger.info(f"df_col:{dataframe.columns}") partner_payment_ids = dataframe.select("partner_payment_id", "effective_payroll_date").distinct() self.logger.info(f"head:{partner_payment_ids.head(3)}") self.logger.info(f"count:{partner_payment_ids.count()}") self.logger.info(f"c:{partner_payment_ids.columns}") temp_table_name = self._create_and_load_table(partner_payment_ids) table_schema = "partner_staging" update_query = """UPDATE abc.partner_payment as pp JOIN partner_staging.{temp_table_name} as t on pp.id = t.partner_payment_id SET pp.reconciled = 1, pp.modify_ts = NOW() , pp.modify_user = 'ok successfull'""" self.db.execute_insert_update(query_sql=update_query) self.db.drop_table(table_schema, temp_table_name) self.logger.info("End reconcile_partner_payments") def execute_insert_update(self, query_sql, args=()) -> None: #method being used in second_function try: cursor = self.get_cursor() self.start_transaction() cursor.execute(query_sql, args) self._connection.handle_unread_result() self.commit() except Error as e: stack = traceback.format_exc() self.rollback() print(stack) logger_func = Cache()[keys.LOGGER].info if keys.LOGGER in Cache() else print logger_func(f"stack_: {stack}") raise e finally: self.close_cursor() def _generate_table_name(self):#generate unique name self.logger.info("Start _generate_table_name") rand_str = random_str(length=4) self.logger.info("End _generate_table_name") return f"temp_{self.partner_name}_{self.company_name}_{self.payroll_name}_{rand_str}"
解决方案
1. 修正SQL占位符未替换的问题
当前update_query中的{temp_table_name}只是字符串模板,未实际替换为生成的临时表名,需添加格式化代码:
update_query = update_query.format(temp_table_name=temp_table_name)
2. 将所有临时表操作纳入同一事务
1412错误的核心原因是临时表的创建/加载与后续更新操作不在同一个事务中,MySQL会判定事务期间表定义发生变化。需把临时表创建、数据加载、更新、删除全部放在同一个事务内执行:
修改reconcile_partner_payments方法:
def reconcile_partner_payments(self, dataframe: DataFrame) -> NoReturn: self.logger.info("Start reconcile_partner_payments") partner_payment_ids = dataframe.select("partner_payment_id", "effective_payroll_date").distinct() temp_table_name = self._generate_table_name() table_schema = "partner_staging" table_with_schema = f"{table_schema}.{temp_table_name}" # 定义事务内的所有操作逻辑 def transaction_ops(cursor): # 1. 显式创建临时表,指定结构避免自动推断的不确定性 create_sql = f""" CREATE TABLE {table_with_schema} ( partner_payment_id INT, effective_payroll_date DATE ) ENGINE=InnoDB """ cursor.execute(create_sql) # 2. 批量插入数据到临时表(替代JDBC跨连接写入) data_rows = partner_payment_ids.rdd.map(lambda row: (row.partner_payment_id, row.effective_payroll_date)).collect() insert_sql = f"INSERT INTO {table_with_schema} VALUES (%s, %s)" cursor.executemany(insert_sql, data_rows) # 3. 执行更新操作 update_sql = f""" UPDATE abc.partner_payment as pp JOIN {table_with_schema} as t ON pp.id = t.partner_payment_id SET pp.reconciled = 1, pp.modify_ts = NOW(), pp.modify_user = 'ok successfull' """ cursor.execute(update_sql) # 4. 删除临时表 drop_sql = f"DROP TABLE {table_with_schema}" cursor.execute(drop_sql) # 调用事务执行方法 self.db.execute_in_transaction(transaction_ops) self.logger.info("End reconcile_partner_payments")
新增事务执行方法替代原execute_insert_update:
def execute_in_transaction(self, transaction_func) -> None: try: cursor = self.get_cursor() self.start_transaction() transaction_func(cursor) self._connection.handle_unread_result() self.commit() except Error as e: stack = traceback.format_exc() self.rollback() print(stack) logger_func = Cache()[keys.LOGGER].info if keys.LOGGER in Cache() else print logger_func(f"stack_: {stack}") raise e finally: self.close_cursor()
3. 确保连接会话一致性
临时表是MySQL会话级对象,跨会话无法访问。需保证df_mgr.jdbc_write使用的连接与db对象的连接为同一个会话,若无法复用连接,直接用上述方案中的批量插入方式在当前连接操作。
4. 显式指定临时表结构
避免依赖JDBC自动推断表结构,创建临时表时明确指定字段类型、存储引擎,确保表定义稳定,减少触发1412错误的概率。
内容的提问来源于stack exchange,提问作者vish
相关产品推荐
相关产品推荐

