使用psycopg2批量导入文件时数据重复且日期不匹配问题求助
问题:导入文本文件到PostgreSQL时数据重复且日期匹配错误
我编程经验尚浅,目前尝试用psycopg2将目录下格式统一的文本文件导入PostgreSQL,希望为每条数据添加对应文件的日期列。但导入后出现异常:始终重复最后一个文件的数据,却搭配其他文件的日期。
示例文件内容
file_01/01/2024
VDS_ID,FWY,co 1001,1,LA 1002,1,LA 1003,2,LA 1004,1,LA 1005,1,LA
file_03/01/2024
VDS_ID,FWY,co 1001,1,LA 1002,1,LA 1003,2,LA 1005,1,LA
file_05/01/2024
VDS_ID,FWY,co 1001,1,LA 1002,1,LA 1003,2,LA
导入后异常结果
VDS_ID,FWY,config_date 1001,1,LA,01/01/2024 1002,1,LA,01/01/2024 1003,2,LA,01/01/2024 1001,1,LA,03/01/2024 1002,1,LA,03/01/2024 1003,2,LA,03/01/2024 1001,1,LA,05/01/2024 1002,1,LA,05/01/2024 1003,2,LA,05/01/2024
我的代码
import psycopg2 from psycopg2.extras import execute_values import time import os import subprocess import csv def main(): mypath = r'C:\Downloads\PEMS\Meta' tableName = 'pems_config' temptableName = 'temp_pems_config' # Connect to an existing database (password blanked out question conn = psycopg2.connect(database='Test', user='postgres', password ='****') # Open a cursor to perform database operations cur = conn.cursor() start_time = time.time() print ("Start time: {0}", time.asctime(time.localtime(start_time))) textfiles = [] for file in os.listdir(mypath): if file.endswith(".txt"): ## print(os.path.join(mypath, file)) textfiles.append(file) create_table_sql = """CREATE TABLE pems_config ( "VDS_ID" integer, "Fwy" smallint, "Dir" character(1), "District" smallint, "County" integer, "City" integer, "State_PM" varchar(9), "Abs_PM" real, "Latitude" real, "Longitude" real, "Length" real, "Type" character(2), "Lanes" smallint, "Name" text, "User_ID_1" text, "User_ID_2" text, "User_ID_3" text, "User_ID_4" text, "pems_config_date" date );""" create_temp_table_sql = """CREATE TEMPORARY TABLE temp_pems_config ( "VDS_ID" integer, "Fwy" smallint, "Dir" character(1), "District" smallint, "County" integer, "City" integer, "State_PM" varchar(9), "Abs_PM" real, "Latitude" real, "Longitude" real, "Length" real, "Type" character(2), "Lanes" smallint, "Name" text, "User_ID_1" text, "User_ID_2" text, "User_ID_3" text, "User_ID_4" text );""" drop_table_sql = """DROP TABLE IF EXISTS temp_pems_config""" copy_sql = """COPY temp_pems_config FROM stdin WITH DELIMITER '\t' NULL AS '' csv HEADER""" add_column_sql = """ALTER TABLE temp_pems_config ADD COLUMN pems_config_date date;""" update_date_sql = """UPDATE temp_pems_config SET pems_config_date = %s ;""" append_table_sql = """INSERT INTO pems_config select * FROM temp_pems_config;""" final_table_sql = """CREATE TABLE final_pems_config AS Select "VDS_ID","Dir","District","County","City","State_PM","Abs_PM","Latitude","Longitude","Length","Type", "Lanes","Name","User_ID_1","User_ID_2","User_ID_3","User_ID_4",max("pems_config_date") AS "pems_config_date" from pems_config GROUP BY "VDS_ID","Dir","District","County","City","State_PM","Abs_PM","Latitude","Longitude","Length","Type", "Lanes","Name","User_ID_1","User_ID_2","User_ID_3","User_ID_4";""" cur.execute("SELECT EXISTS (SELECT 1 AS result FROM pg_tables WHERE schemaname = 'public' AND tablename = 'pems_config');") tableExists = cur.fetchone()[0] print ("{0} Exists: {1}".format(tableName, str(tableExists))) if tableExists == False: print ("Creating {0} Table:".format(tableName)) #Execute a command: this creates a new table cur.execute(create_table_sql) print (cur.statusmessage) print ("Time to create {0} Table: {1} s".format(tableName,time.time()-start_time)) for tfile in textfiles: tfile_name = os.path.splitext(tfile)[0] config_date = tfile_name[14:18]+"-"+tfile_name[19:21]+"-"+tfile_name[22:25] ## #check if table exists cur.execute("SELECT EXISTS (SELECT 1 AS result FROM pg_tables WHERE schemaname = 'public' AND tablename = 'temp_pems_config');") tableExists = cur.fetchone()[0] print ("{0} Exists: {1}".format(temptableName, str(tableExists))) ## ## #if table does not exist then create it if tableExists == False: print ("Creating {0} Table:".format(temptableName)) #Execute a command: this creates a new table cur.execute(drop_table_sql) cur.execute(create_temp_table_sql) print (cur.statusmessage) print ("Time to create {0} Table: {1} s".format(tableName,time.time()-start_time)) else: print ("Deleting {0} Table:".format(tableName)) cur.execute(drop_table_sql) print ("Creating {0} Table:".format(tableName)) cur.execute(create_table_sql) print (cur.statusmessage) print ("Time to create {0} Table: {1} s".format(tableName,time.time()-start_time)) with open(os.path.join(mypath,file), 'r') as f: copyDataFile_time = time.time() #copy zfile into CHTV execute_values cur.copy_expert(sql=copy_sql, file=f) cur.execute(add_column_sql) cur.execute(update_date_sql, (config_date,)) cur.execute(append_table_sql) print ("Time to copy {0} Table: {1} mins".format(tableName,(time.time()-copyDataFile_time)/60)) print (cur.statusmessage) cur.execute(final_table_sql) conn.commit() cur.close() Endtime = time.time() print ("Start time: {0}".format(time.asctime(time.localtime(start_time)))) print ("End time: {0} ".format(time.asctime(time.localtime(Endtime)))) print ("Total processing time: {0} Minutes".format((Endtime-start_time)/60)) if __name__ == '__main__': main()
异常原因
- 文件读取变量错误:循环遍历
textfiles时用的是tfile变量,但打开文件时却用了外层遍历目录的file变量,导致每次循环都读取最后一个遍历到的文件,而日期是当前tfile对应的日期,所以所有日期都绑定同一个文件的数据。 - 临时表逻辑混乱:当临时表存在时,错误地删除临时表后重建了主表
pems_config,直接清空之前导入的数据,每次循环重复这个错误操作,进一步导致数据异常。 - 无效代码残留:
execute_values这行代码没有实际作用,属于冗余代码。
解决方法
- 修正文件读取变量:将
open(os.path.join(mypath,file), 'r')改为open(os.path.join(mypath,tfile), 'r'),确保每次循环读取当前迭代的目标文件。 - 简化临时表处理:删除复杂的存在性判断逻辑,每次循环直接删除并重建临时表,避免逻辑冲突:
# 替换原有的临时表判断代码 cur.execute(drop_table_sql) cur.execute(create_temp_table_sql) print(f"重建临时表 {temptableName}") - 移除错误的主表重建代码:删除else分支中重建
pems_config主表的代码,防止已导入数据被反复清空。 - 清理无效代码:删除
execute_values这一行无用代码。 - 优化日期提取逻辑:PostgreSQL的
date类型要求格式为YYYY-MM-DD,建议改用更可靠的方式提取日期,避免索引截取错误:# 假设文件名格式为file_dd/mm/yyyy.txt date_part = tfile_name.split('_')[1] # 提取dd/mm/yyyy部分 day, month, year = date_part.split('/') config_date = f"{year}-{month}-{day}" # 转换为标准日期格式
修正后的核心循环代码
for tfile in textfiles: tfile_name = os.path.splitext(tfile)[0] # 优化日期提取逻辑,适配文件名格式 date_part = tfile_name.split('_')[1] day, month, year = date_part.split('/') config_date = f"{year}-{month}-{day}" # 每次循环重建临时表 cur.execute(drop_table_sql) cur.execute(create_temp_table_sql) print(f"重建临时表 {temptableName}") with open(os.path.join(mypath, tfile), 'r') as f: copyDataFile_time = time.time() cur.copy_expert(sql=copy_sql, file=f) cur.execute(add_column_sql) cur.execute(update_date_sql, (config_date,)) cur.execute(append_table_sql) print(f"导入文件 {tfile} 耗时: {(time.time()-copyDataFile_time)/60} 分钟") print(cur.statusmessage)
内容的提问来源于stack exchange,提问作者aconti74
相关产品推荐
相关产品推荐

