Airflow导入JSON到PostgreSQL时类型不匹配问题求助
解决PostgreSQL导入时的类型不匹配及代码问题
问题根源
date_id字段被错误提取为列表类型(text[]),但PostgreSQL表中该字段定义为date类型,引发类型不匹配错误- 代码存在循环变量未定义、字段提取逻辑错误的问题
分步解决
1. 修正字段提取:将列表转为单个值
检查你的字典列表数据,若date_id是嵌套列表格式(如{"date_id": ["2024-05-20"]}),需提取列表中的单个值:
# 原始字典列表数据示例 raw_data = [{"date_id": ["2024-05-20"], "other_column": "sample_value"}, ...] # 清洗数据,取出date_id的单个值 cleaned_data = [] for item in raw_data: cleaned_item = item.copy() # 处理列表类型的date_id if isinstance(cleaned_item.get("date_id"), list) and cleaned_item["date_id"]: cleaned_item["date_id"] = cleaned_item["date_id"][0] cleaned_data.append(cleaned_item)
2. 统一日期格式适配PostgreSQL
PostgreSQL的date类型仅接受YYYY-MM-DD格式的字符串,若日期格式不符,需转换:
from datetime import datetime for item in cleaned_data: date_str = item["date_id"] if isinstance(date_str, str): try: # 按原始日期格式调整strptime参数,比如原始是"2024/05/20"则用"%Y/%m/%d" date_obj = datetime.strptime(date_str, "%Y/%m/%d") item["date_id"] = date_obj.strftime("%Y-%m-%d") except ValueError: # 处理无效日期,可选择跳过或记录日志 print(f"跳过无效日期数据: {item}") cleaned_data.remove(item)
3. 修复循环变量问题,编写正确的批量插入代码
使用Airflow的PostgresHook简化连接,配合executemany实现高效批量插入:
from airflow.hooks.postgres_hook import PostgresHook # 替换为你的PostgreSQL连接ID pg_hook = PostgresHook(postgres_conn_id="your_postgres_conn") conn = pg_hook.get_conn() cursor = conn.cursor() # 替换为你的表名和对应字段 insert_sql = """ INSERT INTO target_table (date_id, other_column) VALUES (%s, %s) ON CONFLICT DO NOTHING; -- 按需添加冲突处理逻辑 """ # 准备插入的元组列表(确保字段顺序与SQL一致) insert_values = [ (item["date_id"], item["other_column"]) for item in cleaned_data if all(key in item for key in ["date_id", "other_column"]) ] try: cursor.executemany(insert_sql, insert_values) conn.commit() except Exception as e: conn.rollback() raise e finally: cursor.close() conn.close()
4. 从源头优化(Spark阶段)
如果可以,在Spark生成JSON时直接处理date_id,避免后续清洗:
// Spark Scala代码示例 import org.apache.spark.sql.functions.date_format val parquetDf = spark.read.parquet("path/to/input_parquet") val processedDf = parquetDf.withColumn("date_id", date_format(col("date_id"), "yyyy-MM-dd")) processedDf.write.mode("overwrite").json("path/to/output_json")
内容的提问来源于stack exchange,提问作者Shashank Tiwari
相关产品推荐
相关产品推荐

