add workflow 天润Smart-ccc会话数据,dev
This commit is contained in:
parent
84b259a4c0
commit
606ab9c5d4
|
@ -24,7 +24,7 @@ default_args = {
|
|||
}
|
||||
|
||||
dag = DAG('wf_dag_smart_ccc_chat', default_args=default_args,
|
||||
schedule_interval="59 0-23/1 * * *",
|
||||
schedule_interval="0 0-23/1 * * *",
|
||||
catchup=False,
|
||||
dagrun_timeout=timedelta(minutes=160),
|
||||
max_active_runs=3)
|
||||
|
@ -37,7 +37,7 @@ task_failed = EmailOperator (
|
|||
cc=[""],
|
||||
subject="smart_ccc_chat_failed",
|
||||
html_content='<h3>您好,smart_ccc_chat作业失败,请及时处理" </h3>')
|
||||
|
||||
|
||||
chat_records_feign = SSHOperator(
|
||||
ssh_hook=sshHook,
|
||||
task_id='chat_records_feign',
|
||||
|
@ -56,7 +56,7 @@ retries=3,
|
|||
dag=dag)
|
||||
|
||||
chat_records_feign >> chat_records_load
|
||||
|
||||
|
||||
tr_chat_messages_4800 = SSHOperator(
|
||||
ssh_hook=sshHook,
|
||||
task_id='tr_chat_messages_4800',
|
||||
|
@ -65,7 +65,7 @@ params={'my_param':"S98_S_tr_chat_messages"},
|
|||
depends_on_past=False,
|
||||
retries=3,
|
||||
dag=dag)
|
||||
|
||||
|
||||
tr_chat_records_2337 = SSHOperator(
|
||||
ssh_hook=sshHook,
|
||||
task_id='tr_chat_records_2337',
|
||||
|
@ -74,7 +74,7 @@ params={'my_param':"S98_S_tr_chat_records"},
|
|||
depends_on_past=False,
|
||||
retries=3,
|
||||
dag=dag)
|
||||
|
||||
|
||||
t01_ccc_chat_record = SSHOperator(
|
||||
ssh_hook=sshHook,
|
||||
task_id='t01_ccc_chat_record',
|
||||
|
@ -82,7 +82,7 @@ command='/data/airflow/etl/PDM/run_psql.sh {{ ds_nodash }} {{params.my_param}} >
|
|||
params={'my_param':"t01_ccc_chat_record_agi"},
|
||||
depends_on_past=False,
|
||||
retries=3,
|
||||
dag=dag)
|
||||
dag=dag)
|
||||
t01_ccc_chat_message_detail = SSHOperator(
|
||||
ssh_hook=sshHook,
|
||||
task_id='t01_ccc_chat_message_detail',
|
||||
|
@ -90,7 +90,7 @@ command='/data/airflow/etl/PDM/run_psql.sh {{ ds_nodash }} {{params.my_param}} >
|
|||
params={'my_param':"t01_ccc_chat_message_detail_agi"},
|
||||
depends_on_past=False,
|
||||
retries=3,
|
||||
dag=dag)
|
||||
dag=dag)
|
||||
d_ccc_cust_info = SSHOperator(
|
||||
ssh_hook=sshHook,
|
||||
task_id='d_ccc_cust_info',
|
||||
|
@ -98,7 +98,7 @@ command='/data/airflow/etl/COM/run_psql.sh {{ ds_nodash }} {{params.my_param}} >
|
|||
params={'my_param':"d_ccc_cust_info_agi"},
|
||||
depends_on_past=False,
|
||||
retries=3,
|
||||
dag=dag)
|
||||
dag=dag)
|
||||
cust_chat_ccc_record = SSHOperator(
|
||||
ssh_hook=sshHook,
|
||||
task_id='cust_chat_ccc_record',
|
||||
|
@ -106,7 +106,7 @@ command='/data/airflow/etl/COM/run_psql.sh {{ ds_nodash }} {{params.my_param}} >
|
|||
params={'my_param':"cust_chat_ccc_record_agi"},
|
||||
depends_on_past=False,
|
||||
retries=3,
|
||||
dag=dag)
|
||||
dag=dag)
|
||||
cust_chat_record_info = SSHOperator(
|
||||
ssh_hook=sshHook,
|
||||
task_id='cust_chat_record_info',
|
||||
|
@ -114,13 +114,50 @@ command='/data/airflow/etl/MART/run_psql.sh {{ ds_nodash }} {{params.my_param}}
|
|||
params={'my_param':"cust_chat_record_info_agi"},
|
||||
depends_on_past=False,
|
||||
retries=3,
|
||||
dag=dag)
|
||||
chat_records_load >> tr_chat_records_2337
|
||||
chat_records_load >> tr_chat_messages_4800
|
||||
tr_chat_records_2337 >> t01_ccc_chat_record
|
||||
tr_chat_messages_4800 >> t01_ccc_chat_message_detail
|
||||
t01_ccc_chat_record >> d_ccc_cust_info
|
||||
t01_ccc_chat_record >> cust_chat_ccc_record
|
||||
d_ccc_cust_info >> cust_chat_record_info
|
||||
cust_chat_ccc_record >> cust_chat_record_info
|
||||
cust_chat_record_info >> task_failed
|
||||
dag=dag)
|
||||
cust_contact_mapping = SSHOperator(
|
||||
ssh_hook=sshHook,
|
||||
task_id='cust_contact_mapping',
|
||||
command='/data/airflow/etl/COM/run_psql.sh {{ ds_nodash }} {{params.my_param}} >>/data/airflow/logs/run_tpt_{{ds_nodash}}.log 2>&1 ',
|
||||
params={'my_param':"cust_contact_mapping_agi"},
|
||||
depends_on_past=False,
|
||||
retries=3,
|
||||
dag=dag)
|
||||
cust_contact_info = SSHOperator(
|
||||
ssh_hook=sshHook,
|
||||
task_id='cust_contact_info',
|
||||
command='/data/airflow/etl/COM/run_psql.sh {{ ds_nodash }} {{params.my_param}} >>/data/airflow/logs/run_tpt_{{ds_nodash}}.log 2>&1 ',
|
||||
params={'my_param':"cust_contact_info_agi"},
|
||||
depends_on_past=False,
|
||||
retries=3,
|
||||
dag=dag)
|
||||
cust_leads = SSHOperator(
|
||||
ssh_hook=sshHook,
|
||||
task_id='cust_leads',
|
||||
command='/data/airflow/etl/COM/run_psql.sh {{ ds_nodash }} {{params.my_param}} >>/data/airflow/logs/run_tpt_{{ds_nodash}}.log 2>&1 ',
|
||||
params={'my_param':"cust_leads_agi"},
|
||||
depends_on_past=False,
|
||||
retries=3,
|
||||
dag=dag)
|
||||
cust_all_info = SSHOperator(
|
||||
ssh_hook=sshHook,
|
||||
task_id='cust_all_info',
|
||||
command='/data/airflow/etl/MART/run_psql.sh {{ ds_nodash }} {{params.my_param}} >>/data/airflow/logs/run_tpt_{{ds_nodash}}.log 2>&1 ',
|
||||
params={'my_param':"cust_all_info_agi"},
|
||||
depends_on_past=False,
|
||||
retries=3,
|
||||
dag=dag)
|
||||
chat_records_load >> tr_chat_records_2337
|
||||
chat_records_load >> tr_chat_messages_4800
|
||||
tr_chat_records_2337 >> t01_ccc_chat_record
|
||||
tr_chat_messages_4800 >> t01_ccc_chat_message_detail
|
||||
t01_ccc_chat_record >> d_ccc_cust_info
|
||||
t01_ccc_chat_record >> cust_chat_ccc_record
|
||||
d_ccc_cust_info >> cust_chat_record_info
|
||||
cust_chat_ccc_record >> cust_chat_record_info
|
||||
d_ccc_cust_info >> cust_contact_mapping
|
||||
cust_contact_mapping >> cust_contact_info
|
||||
cust_contact_mapping >> cust_leads
|
||||
cust_leads >> cust_all_info
|
||||
cust_contact_info >> cust_all_info
|
||||
cust_all_info >> task_failed
|
||||
|
|
Loading…
Reference in New Issue