Dagconfig

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":[[]]}
# 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 succeeded
  • all_failed: All upstream failed
  • all_done: All upstream completed
  • one_failed: At least one upstream failed
  • one_success: At least one upstream succeeded
  • none_failed: No upstream failed
  • dummy: No dependency