from airflow.sdk import dag
from airflow.sdk.bases.operator import chain
from iqvia_nlp_provider.hooks.legacy_webhook import ResultOptions
from iqvia_nlp_provider.operators.i2e import (
I2EDeleteResourcesOperator,
I2EMakeIndexOperator,
I2EPrintTaskLogOperator,
I2ERunQueryOperator,
)
from iqvia_nlp_provider.transfers import Webhook
from iqvia_nlp_provider.transfers.s3 import I2EToS3Operator, S3ToI2EOperator
@dag(dag_id="example_dags.i2e_legacy_webhooks", schedule=None, catchup=False)
def i2e_legacy_webhooks():
"""
DAG showing how to interact with legacy NLP Data Factory webhooks.
This DAG calls a before hook (via the S3ToI2EOperator) and an after hook (via the I2EToS3Operator).
Requirements to use this DAG:
- An I2E connection named `i2e_default`.
- An AWS connection named `aws_conn_id`.
- An Legacy NLP Data Factory Webhook connection named `webhook_default`. The webhook server must
have a '/before-hook/' endpoint and an '/after-hook/' endpoint.
- An S3 bucket named `s3-bucket`.
- A source data file in the S3 bucket with the key `source.txt`.
- 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.
"""
s3_to_i2e_task = i2e_source_data_uri = S3ToI2EOperator(
aws_conn_id="aws_conn_id",
i2e_conn_id="i2e_default",
aws_bucket_name="s3-bucket",
aws_key="source.txt",
legacy_before_hook=Webhook(
conn_id="webhook_default",
endpoint="before-hook/",
),
)
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 = i2e_query_results_uri = I2ERunQueryOperator(
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,
"queryProperties": {"outputSettings": {"outputMode": "html"}},
},
)
i2e_to_s3_task = I2EToS3Operator(
i2e_conn_id="i2e_default",
aws_conn_id="aws_conn_id",
i2e_query_results_uri=i2e_query_results_uri.output,
aws_bucket_name="s3-bucket",
legacy_after_hook=Webhook(
conn_id="webhook_default",
endpoint="after-hook/",
result_options=ResultOptions(
query_results=True,
cached_documents=True,
highlight=True,
highlighted_documents="content_and_stylesheet",
),
),
i2e_source_data_uri=i2e_source_data_uri.output,
)
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.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 >> i2e_print_task_log_task
chain(
s3_to_i2e_task,
i2e_make_index_task,
i2e_run_query_task,
i2e_to_s3_task,
i2e_delete_resources_task,
)
i2e_legacy_webhooks()