i2e_query_with_multiple_queries.pyΒΆ

from airflow.sdk import dag
from airflow.sdk.bases.operator import chain

from iqvia_nlp_provider.operators.i2e import (
    I2EDeleteResourcesOperator,
    I2EMakeIndexOperator,
    I2EPrintTaskLogOperator,
    I2ERunQueryOperator,
)
from iqvia_nlp_provider.transfers.local import I2EToLocalFilesystemOperator, LocalFilesystemToI2EOperator


@dag(dag_id="example_dags.i2e_query_with_multiple_queries", schedule=None, catchup=False)
def i2e_query_with_multiple_queries():
    """
    DAG showing how to use I2E operators to upload source data from a local file to I2E,
    index the source data, run two queries against the index, write the query results
    to local files, and delete all transient resources from I2E.

    Requirements to use this DAG:
        - An I2E connection named `i2e_default`.
        - A source file at '/work_dir/source.txt' on the Airflow worker's filesystem.
        - A directory named '/work_dir/results/' on the Airflow worker's filesystem.
        - An index template named `data-factory-workflow-example` installed on the I2E server.
        - A query template named `data-factory-workflow-example` installed on the I2E server.
    """
    local_filesystem_to_i2e_task = i2e_source_data_uri = LocalFilesystemToI2EOperator(
        i2e_conn_id="i2e_default",
        local_file_path="/work_dir/source.txt",
    )

    i2e_make_index_task = i2e_index_uri = I2EMakeIndexOperator(
        i2e_conn_id="i2e_default",
        i2e_source_data_uri=i2e_source_data_uri.output,
        i2e_index_template="data-factory-workflow-example",
        i2e_index_settings_overrides={
            "useDaemon": True,
            "checkDaemon": True,
            "ontologyIndexSuppression": "All",
        },
    )

    i2e_print_task_log_task = I2EPrintTaskLogOperator(
        i2e_conn_id="i2e_default",
    )

    i2e_run_query_task_1 = i2e_query_results_uri_1 = I2ERunQueryOperator(
        task_id="i2e_run_query_1",
        i2e_conn_id="i2e_default",
        i2e_index_uri=i2e_index_uri.output,
        i2e_query_template="data-factory-workflow-example",
        i2e_query_settings_overrides={
            "checkClasses": False,
        },
    )

    i2e_run_query_task_2 = i2e_query_results_uri_2 = I2ERunQueryOperator(
        task_id="i2e_run_query_2",
        i2e_conn_id="i2e_default",
        i2e_index_uri=i2e_index_uri.output,
        i2e_query_template="data-factory-workflow-example",
        i2e_query_settings_overrides={
            "checkClasses": False,
        },
    )

    i2e_to_local_filesystem_task_1 = I2EToLocalFilesystemOperator(
        task_id="i2e_to_local_filesystem_1",
        i2e_conn_id="i2e_default",
        i2e_query_results_uri=i2e_query_results_uri_1.output,
        local_folder="/work_dir/results/",
    )

    i2e_to_local_filesystem_task_2 = I2EToLocalFilesystemOperator(
        task_id="i2e_to_local_filesystem_2",
        i2e_conn_id="i2e_default",
        i2e_query_results_uri=i2e_query_results_uri_2.output,
        local_folder="/work_dir/results/",
    )

    i2e_delete_resources_task = I2EDeleteResourcesOperator(
        i2e_conn_id="i2e_default",
        i2e_resource_uris=[
            i2e_source_data_uri.output,
            i2e_index_uri.output,
            i2e_query_results_uri_1.output,
            i2e_query_results_uri_2.output,
        ],
    )

    # Print task log runs only if make index or query fails
    i2e_make_index_task >> i2e_print_task_log_task
    [i2e_run_query_task_1, i2e_run_query_task_2] >> i2e_print_task_log_task

    chain(
        local_filesystem_to_i2e_task,
        i2e_make_index_task,
        [i2e_run_query_task_1, i2e_run_query_task_2],
        [i2e_to_local_filesystem_task_1, i2e_to_local_filesystem_task_2],
        i2e_delete_resources_task,
    )


i2e_query_with_multiple_queries()