【发布时间】:2020-12-09 11:34:39
【问题描述】:
我正在尝试形成一个简单的管道,我首先清理数据框,然后将其存储在 CSV 中。 clean_Data 任务失败,因此任务 store_data 不工作。我在store_data 任务中使用完整路径将数据帧写入为CSV。我还首先尝试在本地读取 clean_Data 中的 csv,然后尝试在不从 CSV 读取的情况下创建数据框,但两者都不起作用。
这是我写的 DAG 脚本:
from airflow import DAG
import os
from datetime import timedelta
from datetime import datetime
from datetime import date
from airflow.operators.bash_operator import BashOperator
from airflow.operators.python_operator import PythonOperator
from airflow.utils.dates import days_ago
# import pandas and numpy
import pandas as pd
import numpy as np
default_args = {
'owner': 'airflow',
'depends_on_past': False,
'start_date': datetime(2020, 12, 5),
'retries': 1
}
def storeData(**context):
df= context['task_instance'].xcom_pull(task_ids='clean_Data')
df.to_csv('/c/Users/HP/airflowhome/dags/data/namesAge2.csv')
def cleanData():
#df = df.read_csv('/c/Users/HP/airflowhome/dags/data/namesAge.csv')
data = {'Name': ['Tom', 'nick', 'krish', 'jack'],
'Age': [np.nan, 21, np.nan, 18]}
df = pd.DataFrame(data)
df = df.fillna(0)
return df
dag = DAG(
'CleaningPipelineDAG',
default_args=default_args,
description='Cleaning Data',
schedule_interval='@once',
)
t1 = PythonOperator(
task_id='clean_Data',
provide_context=True,
python_callable=cleanData,
dag=dag,
)
t2 = PythonOperator(
task_id='store_data',
provide_context=True,
python_callable=storeData,
dag=dag,
)
t1 >> t2
这是日志文件的内容:
[2020-12-09 13:50:55,984] {taskinstance.py:670} INFO - Dependencies all met for <TaskInstance: CleaningPipelineDAG.clean_Data 2020-12-09T11:45:26.969635+00:00 [queued]>
[2020-12-09 13:50:55,992] {taskinstance.py:670} INFO - Dependencies all met for <TaskInstance: CleaningPipelineDAG.clean_Data 2020-12-09T11:45:26.969635+00:00 [queued]>
[2020-12-09 13:50:55,992] {taskinstance.py:880} INFO -
--------------------------------------------------------------------------------
[2020-12-09 13:50:55,992] {taskinstance.py:881} INFO - Starting attempt 2 of 2
[2020-12-09 13:50:55,992] {taskinstance.py:882} INFO -
--------------------------------------------------------------------------------
[2020-12-09 13:50:56,004] {taskinstance.py:901} INFO - Executing <Task(PythonOperator): clean_Data> on 2020-12-09T11:45:26.969635+00:00
[2020-12-09 13:50:56,020] {standard_task_runner.py:54} INFO - Started process 1698 to run task
[2020-12-09 13:50:56,091] {standard_task_runner.py:77} INFO - Running: ['airflow', 'run', 'CleaningPipelineDAG', 'clean_Data', '2020-12-09T11:45:26.969635+00:00', '--job_id', '76', '--pool', 'default_pool', '--raw', '-sd', 'DAGS_FOLDER//c/Users/HP/DAG-3.py', '--cfg_path', '/tmp/tmpj92yjpoi']
[2020-12-09 13:50:56,094] {standard_task_runner.py:78} INFO - Job 76: Subtask clean_Data
[2020-12-09 13:50:56,169] {logging_mixin.py:112} INFO - Running <TaskInstance: CleaningPipelineDAG.clean_Data 2020-12-09T11:45:26.969635+00:00 [running]> on host mg.localdomain
[2020-12-09 13:50:56,189] {taskinstance.py:1150} ERROR - cleanData() got an unexpected keyword argument 'conf'
Traceback (most recent call last):
File "/home/mg/.local/lib/python3.8/site-packages/airflow/models/taskinstance.py", line 984, in _run_raw_task
result = task_copy.execute(context=context)
File "/home/mg/.local/lib/python3.8/site-packages/airflow/operators/python_operator.py", line 113, in execute
return_value = self.execute_callable()
File "/home/mg/.local/lib/python3.8/site-packages/airflow/operators/python_operator.py", line 118, in execute_callable
return self.python_callable(*self.op_args, **self.op_kwargs)
TypeError: cleanData() got an unexpected keyword argument 'conf'
[2020-12-09 13:50:56,191] {taskinstance.py:1187} INFO - Marking task as FAILED. dag_id=CleaningPipelineDAG, task_id=clean_Data, execution_date=20201209T114526, start_date=20201209T115055, end_date=20201209T115056
[2020-12-09 13:51:00,984] {local_task_job.py:102} INFO - Task exited with return code 1
【问题讨论】:
-
究竟是什么不工作?例外?有什么问题?
-
@Elad 我已经用日志文件更新了我的问题。在 Airflow UI 本身上,clean_Data 任务的状态为失败。
标签: python-3.x pandas airflow