Introduction #
DAG (Directed Acyclic Graph) is the core concept of SmartPip, used for defining task orchestration in data pipelines. Through DAG configuration, you can flexibly define execution order, dependencies, and trigger rules between tasks.
DAG Syntax Quick Reference #
# Task definition
#driver_name task_name params --comment
# DAG separator (4+ slashes)
////
# Dependencies
taskA >> taskB >> taskC # Serial dependency
[taskA, taskB] >> taskC # Parallel dependency
Custom Parameters #
-- Current time
report_time = datetime.datetime.now()
-- Yesterday
report_time = datetime.datetime.now() - datetime.timedelta(days=1)
-- Last day of previous month
report_time = datetime.datetime.now().replace(day=1) - datetime.timedelta(days=1)
-- Format string
ZYMD = report_time.strftime('%Y%m%d')
-- f-string with variables
PARAM = f"where create_time >= '{ZYMD}' "
Configuration Method #
-- Each JOB is a node with driver name and task name
#starrocks_sql sql_filename
-- With params and comment
#impala_sql sql_filename ZYMD,MSG -- comment
-- DAG separator: 4+ slashes after JOB definitions
////
-- Dependencies
sql_filename1 >> sql_filename2 >> sql_filename3
Supported Drivers #
| Driver | Description |
|---|---|
| datax | DataX extraction |
| kafkastarrocks | Stream load Kafka to StarRocks |
| apistarrocks | API to StarRocks |
| starrocks_sql | Execute SQL on StarRocks |
| flinkcdc | Real-time sync with auto table creation |
| hive_sql | Execute Hive SQL |
| impala_sql | Execute Impala SQL |
| oracle_sql | Execute Oracle SQL |
| gp_sql | Execute Greenplum/PostgreSQL SQL |
| mssql_sql | Execute SQL Server SQL |
| myql_sql | Execute MySQL SQL |
| sp | Execute stored procedure |
| py | Execute Python script |
| diy | Custom task |
| dataset | Query/validate any database |
| refresh_quality | Refresh data quality |
| refresh_smc | Refresh SmartChart dashboard |
| link | DAG dependency |
| branch | Conditional branch execution |
| trigger | Trigger target DAG |
| validate | Data validation via SELECT SQL |
| ktr | Execute Kettle transform |
| kjb | Execute Kettle job |
dataset — Query & Validate #
# Usage: #dataset jobname id remark[optional] maillist[optional]
# remark options:
# (empty) - Error if no results
# info - Print results only
# e1 - Error if results exist
# e2 - Email notification of results
# e3 - Retry every 15 min for 2 hours if no results
# Advanced: get_dataset(id, param=None) returns {"result":"success","data":[[]]}
link — DAG Dependency #
# Usage: #link DAG_name sleeptime[optional] maxtime[optional]
#link DAG_A -- If DAG_A fails or is running, subsequent tasks won't execute
#link DAG_A 60 3600 -- Poll every 60s, error after 3600s timeout
branch / trigger #
def inv_validate():
dataset = get_dataset(407)
if len(dataset['data']) > 1:
return 'task1'
return 'task2'
#impala_sql task1
#trigger task2 targetdagname
#branch branch_job inv_validate
/////////////////////////
branch_job >> [task1, task2]
diy — Custom Task #
def refresh_XX():
import requests
response = requests.get('', verify=False).json()
if response['status'] != 200:
raise Exception(str(response))
#diy refresh_job refresh_XX -- Note: jobname must differ from function name
Trigger Rules #
all_success: All upstream succeededall_failed: All upstream failedall_done: All upstream completedone_failed: At least one upstream failedone_success: At least one upstream succeedednone_failed: No upstream faileddummy: No dependency