如何基于key_hash列在Pandas中用read_sql实现SQL查询
基于key_hash的SQL查询与DataFrame更新操作指南
1. 从SQL库查询对应key_hash的历史数据
先提取当前DataFrame中的key_hash列表,用参数化查询构造SQL语句(避免SQL注入风险),执行后获取历史数据:
import pandas as pd # 假设已建立数据库连接conn,当前DataFrame名为new_df key_hashes = new_df['key_hash'].tolist() # 构造参数化查询(占位符根据数据库类型调整:MySQL/PostgreSQL用%s,SQLite用?) query = """ SELECT * FROM your_table_name WHERE key_hash IN ({placeholders}) """ placeholders = ', '.join(['%s'] * len(key_hashes)) final_query = query.format(placeholders=placeholders) # 执行查询获取历史数据 historical_df = pd.read_sql(final_query, conn, params=key_hashes)
替换your_table_name为实际的SQL表名
2. 合并新旧数据到同一DataFrame
通过左连接将历史数据与当前DataFrame合并,保留所有当前数据行,同时匹配对应历史记录:
# suffixes用于区分新旧字段,避免重名 merged_df = new_df.merge(historical_df, on='key_hash', how='left', suffixes=('_current', '_history'))
3. 执行阈值比较等运算
基于合并后的DataFrame完成指标计算,示例如下:
# 计算number_1当前值与历史值的差值(历史值为空时用0填充) merged_df['number_1_diff'] = merged_df['number_1_current'] - merged_df['number_1_history'].fillna(0) # 判断float_1是否超过阈值0.5 merged_df['float_1_over_threshold'] = merged_df['float_1_current'] > 0.5
4. 将更新后的数据同步回SQL库
根据数据量和数据库类型选择合适的更新方式:
方式一:删旧插新(适合小数据量)
先删除SQL库中已存在的key_hash记录,再插入更新后的数据:
# 删除已有记录 delete_query = """ DELETE FROM your_table_name WHERE key_hash IN ({placeholders}) """ final_delete_query = delete_query.format(placeholders=placeholders) cursor = conn.cursor() cursor.execute(final_delete_query, key_hashes) conn.commit() # 插入更新后的数据 merged_df.to_sql('your_table_name', conn, if_exists='append', index=False)
方式二:原生Upsert(高效推荐)
利用数据库原生的Upsert语法(不同数据库语法不同),实现存在则更新、不存在则插入:
PostgreSQL示例(ON CONFLICT)
upsert_query = """ INSERT INTO your_table_name (key_hash, hash_1, hash_2, date_updated, date_created, number_1, number_2, float_1, float_2, number_1_diff, float_1_over_threshold) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s) ON CONFLICT (key_hash) DO UPDATE SET hash_1 = EXCLUDED.hash_1, date_updated = EXCLUDED.date_updated, number_1 = EXCLUDED.number_1, number_1_diff = EXCLUDED.number_1_diff, float_1_over_threshold = EXCLUDED.float_1_over_threshold """ # 整理待插入数据 data_to_upsert = merged_df[['key_hash', 'hash_1_current', 'hash_2_current', 'date_updated_current', 'date_created_current', 'number_1_current', 'number_2_current', 'float_1_current', 'float_2_current', 'number_1_diff', 'float_1_over_threshold']].values.tolist() # 批量执行Upsert cursor = conn.cursor() cursor.executemany(upsert_query, data_to_upsert) conn.commit()
MySQL用ON DUPLICATE KEY UPDATE,SQL Server用MERGE,需根据实际数据库调整语法
内容的提问来源于stack exchange,提问作者johnnyb
相关产品推荐
相关产品推荐

