|
7 | 7 | from tempfile import TemporaryDirectory |
8 | 8 |
|
9 | 9 | from airflow import DAG, configuration |
10 | | -from airflow.operators.dummy import DummyOperator |
| 10 | +from airflow.operators.empty import EmptyOperator |
11 | 11 | from airflow.operators.python import PythonOperator |
12 | 12 |
|
13 | 13 | from polygonetl.cli import ( |
|
34 | 34 |
|
35 | 35 |
|
36 | 36 | def build_export_dag( |
37 | | - dag_id, |
38 | | - provider_uris, |
39 | | - provider_uris_archival, |
40 | | - output_bucket, |
41 | | - export_start_date, |
42 | | - export_end_date=None, |
43 | | - notification_emails=None, |
44 | | - export_schedule_interval='0 0 * * *', |
45 | | - export_max_workers=10, |
46 | | - export_traces_max_workers=10, |
47 | | - export_batch_size=200, |
48 | | - export_max_active_runs=None, |
49 | | - export_max_active_tasks=None, |
50 | | - export_retries=5, |
51 | | - **kwargs |
| 37 | + dag_id, |
| 38 | + provider_uris, |
| 39 | + provider_uris_archival, |
| 40 | + output_bucket, |
| 41 | + export_start_date, |
| 42 | + export_end_date=None, |
| 43 | + notification_emails=None, |
| 44 | + export_schedule='0 0 * * *', |
| 45 | + export_max_workers=10, |
| 46 | + export_traces_max_workers=10, |
| 47 | + export_batch_size=200, |
| 48 | + export_max_active_runs=None, |
| 49 | + export_max_active_tasks=None, |
| 50 | + export_retries=5, |
| 51 | + **kwargs |
52 | 52 | ): |
53 | 53 | default_dag_args = { |
54 | 54 | "depends_on_past": False, |
@@ -82,7 +82,7 @@ def build_export_dag( |
82 | 82 |
|
83 | 83 | dag = DAG( |
84 | 84 | dag_id, |
85 | | - schedule_interval=export_schedule_interval, |
| 85 | + schedule=export_schedule, |
86 | 86 | default_args=default_dag_args, |
87 | 87 | max_active_runs=export_max_active_runs, |
88 | 88 | max_active_tasks=export_max_active_tasks, |
@@ -345,7 +345,7 @@ def add_export_task( |
345 | 345 | return None |
346 | 346 |
|
347 | 347 | # Operators |
348 | | - export_complete = DummyOperator(task_id="export_complete", dag=dag) |
| 348 | + export_complete = EmptyOperator(task_id="export_complete", dag=dag) |
349 | 349 |
|
350 | 350 | export_blocks_and_transactions_operator = add_export_task( |
351 | 351 | export_blocks_and_transactions_toggle, |
|
0 commit comments