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

Airflow双任务执行报错求助:JSON解析失败问题

Airflow任务JSON解析错误排查与修复

问题背景

使用Airflow实现两个定时任务:

  • 任务1:每分钟调用OpenWeatherMap API,将多城市天气数据保存为JSON文件
  • 任务2:读取最近20个JSON文件,提取temperature、city、pression字段生成CSV文件

任务依赖已设置为task1>>task2,但运行时任务2触发JSONDecodeError,错误提示为Expecting value: line 1 column 1 (char 0),怀疑代码存在问题,考虑使用XCom传递信息。

原代码

import requests
import json
import datetime
import os
import pandas as pd
from airflow import DAG
from airflow.operators.bash import BashOperator
from airflow.operators.python import PythonOperator
from airflow.utils.dates import days_ago


my_dag = DAG(
    dag_id='eval_airflow',
    description="recup_data and transform data to csv",
    schedule_interval='*/1 * * * *',
    default_args={
        'owner': 'airflow',
        'start_date': days_ago(0),
    }#,
    #catchup=False
)

# 定义从OpenWeatherMap获取数据的函数
def recup_data():
    filepath = '/app/raw_files'
    # 创建存储请求结果的目录'/app/raw_files'
    if os.path.exists(filepath) == False:
            os.makedirs(filepath, mode = 511, exist_ok= True)
    # 切换到'/app/raw_files'目录
    os.chdir(filepath)
    # 创建需要获取天气数据的城市列表
    villes = ["paris", "london","washington"]
    cities = {}
    for ville in villes:
            r = requests.get(f"https://api.openweathermap.org/data/2.5/weather?q={ville}&appid=0eb6409c528ceeabc733ad3b07a67b58")
            cities[ville] = r.json()
# 获取当前时间
    now = datetime.datetime.now()
# 基于时间生成文件名
    filename = f"{now.year}-{now.month}-{now.day} {now.hour}:{now.minute}.json"
# 写入文件
    with open(filename, 'w') as file:
    # 将数据以JSON格式写入文件
        json.dump(cities, file)
        r.status_code
    return r.status_code
    


def transform_data_into_csv(n_files=None, filename='data.csv'):
    parent_folder = '/app/raw_files'
    files = sorted(os.listdir(parent_folder), reverse=True)[:20]
    if n_files:
        files = files[:n_files]
    dfs = []
    print('dfs', dfs)
    for f in files:
        with open(os.path.join(parent_folder, f), 'r') as file:
            data_temp = json.load(file)
        for data_city in data_temp:
            dfs.append(
               {
                    'temperature': data_temp[data_city]['main']['temp'],
                    'city': data_temp [data_city]['name'],
                    'pression': data_temp[data_city]['main']['pressure'],
                    'date': f.split('.')[0]
                }
            )
    df = pd.DataFrame(dfs)
    df.to_csv(os.path.join('/app/clean_data', filename), index=False)


task1 = PythonOperator(
    task_id='task1_recup_data',
    python_callable=recup_data,
    dag=my_dag
)

task2 = PythonOperator(
    task_id='task2_Transform_data_into_csv',
    python_callable=transform_data_into_csv,
    dag=my_dag
)

task1 >> task2

错误日志

File "/opt/airflow/dags/tache1.py", line 58, in transform_data_into_csv
    data_temp = json.load(file)
  File "/usr/local/lib/python3.6/json/__init__.py", line 299, in load
    parse_constant=parse_constant, object_pairs_hook=object_pairs_hook, **kw)
  File "/usr/local/lib/python3.6/json/__init__.py", line 354, in loads
    return _default_decoder.decode(s)
  File "/usr/local/lib/python3.6/json/decoder.py", line 339, in decode
    obj, end = self.raw_decode(s, idx=_w(s, 0).end())
  File "/usr/local/lib/python3.6/json/decoder.py", line 357, in raw_decode
    raise JSONDecodeError("Expecting value", s, err.value) from None
json.decoder.JSONDecodeError: Expecting value: line 1 column 1 (char 0)
[2023-01-12 16:52:03,300] {taskinstance.py:1551} INFO - Marking task as FAILED. dag_id=eval_***, task_id=task2_Transform_data_into_csv, execution_date=20230112T165100, start_date=20230112T165203, end_date=20230112T165203
[2023-01-12 16:52:03,334] {local_task_job.py:151} INFO - Task exited with return code 1

问题排查与修复

1. 核心错误原因

JSONDecodeError: Expecting value: line 1 column 1 表示读取的文件是空的或不是有效的JSON格式,可能触发点:

  • 任务1生成的JSON文件为空(API请求失败但未处理)
  • 任务2读取到了非JSON文件(比如临时文件、空文件)
  • 文件路径或权限问题导致无法正确写入/读取

2. 针对性修复方案

(1)任务1:增强API请求容错与文件写入校验

修改recup_data函数,添加请求状态码校验,确保只有成功请求的数据才写入文件,同时避免生成空文件:

def recup_data(**context):
    filepath = '/app/raw_files'
    # 简化目录创建写法
    os.makedirs(filepath, mode=0o777, exist_ok=True)
    villes = ["paris", "london","washington"]
    cities = {}
    for ville in villes:
        r = requests.get(f"https://api.openweathermap.org/data/2.5/weather?q={ville}&appid=0eb6409c528ceeabc733ad3b07a67b58")
        # 仅当请求成功时才保存数据
        if r.status_code == 200:
            cities[ville] = r.json()
        else:
            print(f"获取{ville}数据失败,状态码: {r.status_code}")
    
    # 只有获取到有效数据时才生成文件
    if cities:
        now = datetime.datetime.now()
        # 文件名替换空格为下划线,避免跨平台问题
        filename = f"{now.year}-{now.month}-{now.day}_{now.hour}-{now.minute}.json"
        file_path = os.path.join(filepath, filename)
        with open(file_path, 'w') as file:
            json.dump(cities, file)
        # 用XCom传递生成的文件名,让任务2精准读取最新文件(可选)
        context['ti'].xcom_push(key='generated_file', value=filename)
    return len(cities)  # 返回成功获取的城市数量

(2)任务2:添加文件有效性校验与路径处理

修改transform_data_into_csv函数,跳过空文件和非JSON文件,同时确保clean_data目录存在:

def transform_data_into_csv(**context):
    parent_folder = '/app/raw_files'
    # 确保目标目录存在
    os.makedirs('/app/clean_data', exist_ok=True)
    
    # 过滤出所有JSON文件
    all_files = [f for f in os.listdir(parent_folder) if f.endswith('.json')]
    # 按文件修改时间排序,取最近20个(比文件名排序更可靠)
    files = sorted(all_files, key=lambda x: os.path.getmtime(os.path.join(parent_folder, x)), reverse=True)[:20]
    
    dfs = []
    for f in files:
        file_path = os.path.join(parent_folder, f)
        # 跳过空文件
        if os.path.getsize(file_path) == 0:
            print(f"跳过空文件: {f}")
            continue
        try:
            with open(file_path, 'r') as file:
                data_temp = json.load(file)
            # 校验数据结构是否合法
            for city_key, city_data in data_temp.items():
                if 'main' in city_data and 'name' in city_data:
                    dfs.append({
                        'temperature': city_data['main']['temp'],
                        'city': city_data['name'],
                        'pression': city_data['main']['pressure'],
                        'date': f.split('.')[0].replace('_', ' ')
                    })
                else:
                    print(f"文件{f}中城市{city_key}的数据结构不合法")
        except json.JSONDecodeError:
            print(f"解析JSON文件失败: {f}")
            continue
    
    if dfs:
        df = pd.DataFrame(dfs)
        df.to_csv(os.path.join('/app/clean_data', 'data.csv'), index=False)
    else:
        print("无有效数据可生成CSV")

(3)Airflow任务配置优化

更新任务定义,添加provide_context=True以支持XCom(如果使用):

task1 = PythonOperator(
    task_id='task1_recup_data',
    python_callable=recup_data,
    provide_context=True,
    dag=my_dag
)

task2 = PythonOperator(
    task_id='task2_Transform_data_into_csv',
    python_callable=transform_data_into_csv,
    provide_context=True,
    dag=my_dag
)

task1 >> task2

3. 可选优化:使用XCom精准传递文件信息

如果不需要处理历史文件,仅处理任务1刚生成的文件,可通过XCom让任务2直接读取任务1生成的文件名,避免遍历所有文件:

def transform_data_into_csv(**context):
    # 获取任务1传递的文件名
    generated_file = context['ti'].xcom_pull(task_ids='task1_recup_data', key='generated_file')
    if not generated_file:
        print("任务1未生成有效文件")
        return
    
    parent_folder = '/app/raw_files'
    file_path = os.path.join(parent_folder, generated_file)
    # 后续处理逻辑同上,仅处理该文件
    # ...

内容的提问来源于stack exchange,提问作者yannick icard

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 21:30:44