You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何每日自动将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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.02 11:11:29