Airflow中SqlAlchemy报update parameter必须为非空字典错误求助
问题解决:SQLAlchemy 1.3 下 on_duplicate_key_update 报错
问题根源是SQLAlchemy版本差异:1.4支持直接用statement.inserted作为on_duplicate_key_update的参数,但1.3.24不支持这种简化写法,必须显式传递非空的字段更新映射字典。
修改方案
把statement.on_duplicate_key_update(statement.inserted)替换为手动构造更新规则,用inserted别名引用插入的字段值,具体代码如下:
def process_item(self, item, spider): data = ItemAdapter(item) print(data) with self.engine.connect() as conn: insert_stmt = insert(self.db_table).values(**data) # 构造更新字典:每个字段对应插入语句中的值 update_dict = {col: insert_stmt.inserted[col] for col in data.keys()} # 执行带重复键更新的插入操作 stmt_with_duplicate = insert_stmt.on_duplicate_key_update(**update_dict) conn.execute(stmt_with_duplicate) return item
补充优化(可选)
如果self.db_table是SQLAlchemy的Table对象,直接遍历表的列名构造更新字典更安全,可避免Item字段与表字段不匹配的问题:
update_dict = {col.name: insert_stmt.inserted[col.name] for col in self.db_table.columns}
这样修改后,就能兼容SQLAlchemy 1.3.24版本的语法要求,解决ValueError: update parameter must be a non-empty dictionary的报错。
内容的提问来源于stack exchange,提问作者Arthur Caldas
相关产品推荐
相关产品推荐

