批量向PostgreSQL插入数据时timestamp与bigint类型不匹配错误的排查求助
大家好,我最近在做批量向PostgreSQL插入数据的脚本时遇到了类型不匹配的问题,折腾了好一阵都没解决,想请各位帮忙看看问题出在哪。
先给大家梳理下我的脚本核心逻辑和步骤:
- 实例化
PostgreSQLDB类对象来处理数据库操作 - 通过视图
vw_valid_case_from_db1获取需要保留的case_id列表 - 从db1表提取数据到pandas DataFrame
df_db1_case_table - 用上面的有效case_id过滤
df_db1_case_table得到df_db1_case_table_filtered - 准备把过滤后的DataFrame批量插入到
db2_case_table,两者列名完全匹配
现在遇到的错误如下:
Exception has occurred: DatatypeMismatch column "occurence_timestamp" is of type timestamp without time zone but expression is of type bigint
LINE 1: ...7.0, 2259027.0, NULL, 'CA23307772', NULL, '1441', 1689711600...
HINT: You will need to rewrite or cast the expression.
File "some_path_to.py", line 170, in
cursor.executemany(insert_query, data_to_insert)
psycopg2.errors.DatatypeMismatch: column "occurence_timestamp" is of type timestamp without time zone but expression is of type bigint
LINE 1: ...7.0, 2259027.0, NULL, 'CA23307772', NULL, '1441', 1689711600...
HINT: You will need to rewrite or cast the expression.
我尝试过处理这些timestamp列:遍历所有时间戳列,判断如果是int64或float64类型就转成datetime格式;如果已经是datetime类型就去掉时区。但执行后还是报同样的错误。
另外我发现,如果把所有时间戳列(['occurence_timestamp', 'reported_timestamp', 'created_timestamp', 'modified_timestamp', 'agency_extract_timestamp', 'city_extract_timestamp', 'pdf_extract_timestamp' ])都从插入列中移除,脚本就能正常运行。
以下是我的完整脚本代码:
import pandas as pd import os import logging from datetime import datetime from helper_db_operation import PostgreSQLDB # Set up logging configuration logging.basicConfig(level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s') # Some database credential code postgres_db = PostgreSQLDB(postgres_user, postgres_password, postgres_host, postgres_db_name) # Query to fetch valid case IDs from db1 view sql_query_get_valid_case_case_from_db1 = """ SELECT case_id FROM vw_valid_case_from_db1 """ # Execute query and load results into a Pandas DataFrame try: df_db1_valid_cases = pd.read_sql(sql_query_get_valid_case_case_from_db1, postgres_db.conn) except Exception as e: logging.error(f"Error while fetching data from db1 view: {e}") raise # 1.2) connect db1_case, then apply the filter from 1.1) to exclude invalid case # Query to fetch all case from db1 table sql_query_get_case_table_from_db1 = """ SELECT * FROM public.db1_case """ # Execute query and load results into a Pandas DataFrame df_db1_case_table = pd.read_sql(sql_query_get_case_table_from_db1, postgres_db.conn) # Apply the filter to include only those case that are in the valid list from df_db1_valid_cases valid_case_ids = df_db1_valid_cases['case_id'].tolist() # Filter df_db1_case_table to include only rows where the ID is in the valid case IDs df_db1_case_table_filtered = df_db1_case_table[df_db1_case_table['id'].isin(valid_case_ids)] # Check the result of the filtering logging.debug(f"Filtered {len(df_db1_case_table_filtered)} valid case.") target_table = 'db2_case_table' # Fetch the target table schema to determine the column names dynamically try: query_table_schema = f""" SELECT column_name FROM information_schema.columns WHERE table_name = '{target_table}'; """ df_table_schema = pd.read_sql(query_table_schema, postgres_db.conn) target_columns = df_table_schema['column_name'].tolist() except Exception as e: logging.error(f"Error fetching schema for table {target_table}: {e}") raise # Convert epoch timestamps to datetime for `occurence_timestamp` and other timestamp columns timestamp_columns = [ 'occurence_timestamp', 'reported_timestamp', 'created_timestamp', 'modified_timestamp', 'agency_extract_timestamp', 'city_extract_timestamp', 'pdf_extract_timestamp' ] for col in timestamp_columns: if col in df_db1_case_table_filtered.columns: if df_db1_case_table_filtered[col].dtype in ['int64', 'float64']: df_db1_case_table_filtered[col] = pd.to_datetime( df_db1_case_table_filtered[col], unit='s', errors='coerce' ) elif pd.api.types.is_datetime64_any_dtype(df_db1_case_table_filtered[col]): df_db1_case_table_filtered[col] = df_db1_case_table_filtered[col].dt.tz_localize(None) # Ensure the correct columns are included in `df_for_insertion` df_for_insertion = df_db1_case_table_filtered[ [col for col in df_db1_case_table_filtered.columns if col in target_columns] ] logging.info(df_for_insertion.dtypes) logging.info(df_for_insertion[['occurence_timestamp']].head()) # Insert the filtered and dynamically mapped data into PostgreSQL try: # Convert the DataFrame into a list of tuples data_to_insert = df_for_insertion.to_records(index=False).tolist() logging.info(data_to_insert[:5]) # Generate the INSERT query dynamically based on the DataFrame columns columns = ', '.join(df_for_insertion.columns) placeholders = ', '.join(['%s'] * len(df_for_insertion.columns)) insert_query = f"INSERT INTO {target_table} ({columns}) VALUES ({placeholders})" # Execute the batch insert with postgres_db.conn.cursor() as cursor: cursor.executemany(insert_query, data_to_insert) postgres_db.conn.commit() except Exception as e: logging.error(f"Error while inserting data into table {target_table}: {e}") raise
附加信息
我打印了DataFrame中相关列的类型,显示确实是datetime格式:
occurence_timestamp datetime64[ns] reported_timestamp datetime64[ns]
同时也输出了occurence_timestamp的样例值:
occurence_timestamp 0 2023-07-18 20:20:00 1 2023-09-21 17:00:00 2 2023-09-21 15:48:00 3 2023-09-21 21:30:00 4 2023-09-11 08:45:00
真心希望大家能帮我找出问题所在,谢谢!
备注:内容来源于stack exchange,提问作者KubiK888

