如何每日自动将CSV数据导入PostgreSQL并实现历史数据查询
自动将每日CSV导入PostgreSQL的实现方案
第一步:准备匹配的数据库表
先创建和CSV字段对应的表,建议新增import_time字段记录导入时间,方便后续追溯历史数据。假设CSV包含location_id(地点ID)、measure_value(测量值)、measure_time(测量时间)三个核心字段,创建表的SQL如下:
CREATE TABLE IF NOT EXISTS location_measures ( id SERIAL PRIMARY KEY, location_id VARCHAR(50) NOT NULL, measure_value NUMERIC(10,2) NOT NULL, measure_time TIMESTAMP NOT NULL, import_time TIMESTAMP DEFAULT CURRENT_TIMESTAMP );
第二步:自动导入的两种实用方法
方法1:Shell脚本+系统定时任务(轻量化首选)
写一个shell脚本处理CSV导入逻辑,再通过crontab定时执行,适合无复杂数据校验的场景。
脚本示例(import_csv.sh):
#!/bin/bash # 配置参数,根据你的实际情况修改 CSV_FILE="/data/daily_measures.csv" DB_NAME="your_database" DB_USER="your_db_user" DB_HOST="localhost" DB_PORT="5432" LOG_FILE="/var/log/csv_import.log" ARCHIVE_DIR="/data/csv_archive" # 检查CSV文件是否存在 if [ ! -f "$CSV_FILE" ]; then echo "$(date +'%Y-%m-%d %H:%M:%S'): 错误:未找到CSV文件 $CSV_FILE" >> $LOG_FILE exit 1 fi # 用psql的\copy命令导入数据(指定字段避免顺序问题) psql -h $DB_HOST -p $DB_PORT -U $DB_USER -d $DB_NAME -c "\copy location_measures(location_id, measure_value, measure_time) FROM '$CSV_FILE' DELIMITER ',' CSV HEADER;" # 记录执行结果并归档文件 if [ $? -eq 0 ]; then echo "$(date +'%Y-%m-%d %H:%M:%S'): CSV数据导入成功" >> $LOG_FILE # 导入后将文件移至归档目录,避免重复导入 mv "$CSV_FILE" "$ARCHIVE_DIR/$(date +'%Y%m%d')_measures.csv" else echo "$(date +'%Y-%m-%d %H:%M:%S'): CSV数据导入失败" >> $LOG_FILE fi
设置定时任务:
执行crontab -e,添加一行配置(比如每天凌晨1点执行导入):
0 1 * * * /bin/bash /path/to/import_csv.sh
别忘了给脚本加执行权限:chmod +x /path/to/import_csv.sh
方法2:Python脚本(支持复杂数据校验)
如果需要对CSV数据做清洗、格式校验再导入,用Python结合psycopg2库更灵活。
脚本示例(import_csv.py):
import os import csv import psycopg2 from datetime import datetime # 配置参数 CSV_PATH = "/data/daily_measures.csv" ARCHIVE_DIR = "/data/csv_archive" LOG_PATH = "/var/log/csv_import.log" DB_CONFIG = { "dbname": "your_database", "user": "your_db_user", "password": "your_db_password", "host": "localhost", "port": "5432" } def log(message): with open(LOG_PATH, "a", encoding="utf-8") as f: f.write(f"{datetime.now().strftime('%Y-%m-%d %H:%M:%S')}: {message}\n") def import_csv(): # 检查文件是否存在 if not os.path.exists(CSV_PATH): log(f"未找到CSV文件:{CSV_PATH}") return # 读取并校验CSV数据 try: with open(CSV_PATH, "r", encoding="utf-8") as f: reader = csv.DictReader(f) rows = list(reader) # 校验必填字段 required_fields = ["location_id", "measure_value", "measure_time"] for idx, row in enumerate(rows): for field in required_fields: if not row.get(field): log(f"第{idx+2}行数据缺失必填字段:{field}") return except Exception as e: log(f"读取CSV文件失败:{str(e)}") return # 连接数据库导入数据 conn = None cur = None try: conn = psycopg2.connect(**DB_CONFIG) cur = conn.cursor() insert_sql = """ INSERT INTO location_measures(location_id, measure_value, measure_time) VALUES (%s, %s, %s) """ for row in rows: cur.execute(insert_sql, ( row["location_id"], float(row["measure_value"]), datetime.strptime(row["measure_time"], "%Y-%m-%d %H:%M:%S") # 匹配你的时间格式 )) conn.commit() # 归档文件 archive_name = f"{datetime.now().strftime('%Y%m%d')}_measures.csv" os.rename(CSV_PATH, os.path.join(ARCHIVE_DIR, archive_name)) log(f"成功导入{len(rows)}条数据") except Exception as e: log(f"数据库导入失败:{str(e)}") if conn: conn.rollback() finally: if cur: cur.close() if conn: conn.close() if __name__ == "__main__": import_csv()
同样用crontab定时执行:
0 1 * * * /usr/bin/python3 /path/to/import_csv.py
第三步:查询历史数据示例
导入完成后,即可通过SQL查询历史数据:
- 查询指定地点的所有历史记录:
SELECT * FROM location_measures WHERE location_id = 'LOC001' ORDER BY measure_time DESC;
- 查询最近7天的所有测量数据:
SELECT * FROM location_measures WHERE measure_time >= NOW() - INTERVAL '7 days';
- 按日期统计各地点的测量平均值:
SELECT location_id, DATE(measure_time) AS measure_date, AVG(measure_value) AS avg_measure FROM location_measures GROUP BY location_id, DATE(measure_time) ORDER BY measure_date DESC;
注意事项
- 确保CSV字段顺序和导入SQL指定的字段顺序一致,避免数据错位
- 数据库用户需要拥有目标表的
INSERT权限,以及CSV文件所在目录的读取权限 - 定期清理日志和归档文件,避免占用过多磁盘空间
- 如果CSV存在重复数据,可给表添加唯一约束(如
UNIQUE(location_id, measure_time)),或在导入脚本中增加去重逻辑
内容的提问来源于stack exchange,提问作者Fatima ezzahra Ait moulay erra
相关产品推荐
相关产品推荐

