Skip to content

Commit d0060e4

Browse files
authored
FLYPIE-315 Refactor indexing dags to commit once at the end (#1810)
1 parent 15405ef commit d0060e4

4 files changed

Lines changed: 111 additions & 12 deletions

File tree

funcake_dags/funcake_dev_index_dag.py

Lines changed: 24 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,11 @@
11
"""DAG to Harvest PA Digital Aggregated OAI-PMH XML & Index to SolrCloud."""
22
from datetime import datetime, timedelta
33
from airflow.sdk import DAG
4+
from airflow.providers.amazon.aws.operators.s3 import S3ListOperator
45
from airflow.providers.standard.operators.empty import EmptyOperator
56
from airflow.providers.standard.operators.python import PythonOperator
67
from airflow.providers.standard.operators.bash import BashOperator
8+
from airflow.providers.http.operators.http import HttpOperator
79
from tulflow import harvest, tasks
810
from airflow.providers.slack.notifications.slack import send_slack_notification
911

@@ -94,12 +96,21 @@
9496
CONFIGSET
9597
)
9698

99+
LIST_INDEX_FILES = S3ListOperator(
100+
task_id="list_index_files",
101+
bucket=AIRFLOW_DATA_BUCKET,
102+
prefix=DAG.dag_id + "/" + TIMESTAMP + "/new-updated/",
103+
delimiter="/",
104+
aws_conn_id="AIRFLOW_S3",
105+
dag=DAG
106+
)
107+
97108
COMBINE_INDEX = BashOperator(
98109
task_id="combine_index",
99110
bash_command=FUNCAKE_INDEX_BASH,
100111
env={
101112
"BUCKET": AIRFLOW_DATA_BUCKET,
102-
"FOLDER": DAG.dag_id + "/" + TIMESTAMP + "/new-updated/",
113+
"DATA": "{{ ti.xcom_pull(task_ids='list_index_files') | tojson }}",
103114
"INDEXER": "funnel_cake_index",
104115
"SOLR_URL": SOLR_COLL_ENDPT,
105116
"SOLR_AUTH_USER": "{{ conn.get('SOLRCLOUD-WRITER').login or '' }}",
@@ -113,13 +124,23 @@
113124
dag=DAG
114125
)
115126

127+
SOLR_COMMIT = HttpOperator(
128+
task_id="solr_commit",
129+
method="GET",
130+
http_conn_id=SOLR_CONN_ID,
131+
endpoint="/solr/" + COLLECTION + "/update?commit=true",
132+
dag=DAG
133+
)
134+
116135
SOLR_ALIAS_SWAP = tasks.swap_sc_alias(DAG, SOLR_CONN_ID, COLLECTION, ALIAS)
117136
SUCCESS = EmptyOperator(
118137
task_id='success',
119138
on_success_callback=[slackpostonsuccess])
120139

121140
# SET UP TASK DEPENDENCIES
122141
CREATE_COLLECTION.set_upstream(HARVEST_OAI)
123-
COMBINE_INDEX.set_upstream(CREATE_COLLECTION)
124-
SOLR_ALIAS_SWAP.set_upstream(COMBINE_INDEX)
142+
LIST_INDEX_FILES.set_upstream(CREATE_COLLECTION)
143+
COMBINE_INDEX.set_upstream(LIST_INDEX_FILES)
144+
SOLR_COMMIT.set_upstream(COMBINE_INDEX)
145+
SOLR_ALIAS_SWAP.set_upstream(SOLR_COMMIT)
125146
SUCCESS.set_upstream(SOLR_ALIAS_SWAP)

funcake_dags/funcake_prod_index_dag.py

Lines changed: 24 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,11 @@
11
"""DAG to Harvest PA Digital Aggregated OAI-PMH XML & Index to SolrCloud."""
22
from datetime import datetime, timedelta
33
from airflow.sdk import DAG
4+
from airflow.providers.amazon.aws.operators.s3 import S3ListOperator
45
from airflow.providers.standard.operators.empty import EmptyOperator
56
from airflow.providers.standard.operators.python import PythonOperator
67
from airflow.providers.standard.operators.bash import BashOperator
8+
from airflow.providers.http.operators.http import HttpOperator
79
from tulflow import harvest, tasks
810
from airflow.providers.slack.notifications.slack import send_slack_notification
911

@@ -93,12 +95,21 @@
9395
CONFIGSET
9496
)
9597

98+
LIST_INDEX_FILES = S3ListOperator(
99+
task_id="list_index_files",
100+
bucket=AIRFLOW_DATA_BUCKET,
101+
prefix=DAG.dag_id + "/" + TIMESTAMP + "/new-updated/",
102+
delimiter="/",
103+
aws_conn_id="AIRFLOW_S3",
104+
dag=DAG
105+
)
106+
96107
COMBINE_INDEX = BashOperator(
97108
task_id="combine_index",
98109
bash_command=FUNCAKE_INDEX_BASH,
99110
env={
100111
"BUCKET": AIRFLOW_DATA_BUCKET,
101-
"FOLDER": DAG.dag_id + "/" + TIMESTAMP + "/new-updated/",
112+
"DATA": "{{ ti.xcom_pull(task_ids='list_index_files') | tojson }}",
102113
"INDEXER": "funnel_cake_index",
103114
"SOLR_URL": SOLR_COLL_ENDPT,
104115
"SOLR_AUTH_USER": "{{ conn.get('SOLRCLOUD-WRITER').login or '' }}",
@@ -112,13 +123,23 @@
112123
dag=DAG
113124
)
114125

126+
SOLR_COMMIT = HttpOperator(
127+
task_id="solr_commit",
128+
method="GET",
129+
http_conn_id=SOLR_CONN_ID,
130+
endpoint="/solr/" + COLLECTION + "/update?commit=true",
131+
dag=DAG
132+
)
133+
115134
SOLR_ALIAS_SWAP = tasks.swap_sc_alias(DAG, SOLR_CONN_ID, COLLECTION, ALIAS)
116135
SUCCESS = EmptyOperator(
117136
task_id='success',
118137
on_success_callback=[slackpostonsuccess])
119138

120139
# SET UP TASK DEPENDENCIES
121140
CREATE_COLLECTION.set_upstream(HARVEST_OAI)
122-
COMBINE_INDEX.set_upstream(CREATE_COLLECTION)
123-
SOLR_ALIAS_SWAP.set_upstream(COMBINE_INDEX)
141+
LIST_INDEX_FILES.set_upstream(CREATE_COLLECTION)
142+
COMBINE_INDEX.set_upstream(LIST_INDEX_FILES)
143+
SOLR_COMMIT.set_upstream(COMBINE_INDEX)
144+
SOLR_ALIAS_SWAP.set_upstream(SOLR_COMMIT)
124145
SUCCESS.set_upstream(SOLR_ALIAS_SWAP)

funcake_dags/scripts/index.sh

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,9 +19,30 @@ gem install bundler
1919
bundle config set force_ruby_platform true
2020
bundle install
2121

22+
# Disable Traject's automatic commit on indexer shutdown so Airflow can
23+
# perform one explicit Solr commit after all indexing tasks finish.
24+
sed -i.bak 's/"solr_writer.commit_on_close": true/"solr_writer.commit_on_close": false/' lib/$INDEXER.rb
25+
rm lib/$INDEXER.rb.bak
26+
27+
# grab list of items from designated aws bucket (creds are envvars), then index each item
28+
if [ -n "$DATA" ]; then
29+
RESP=$(echo "$DATA" | jq -r '.[]')
30+
else
31+
RESP=`aws s3 ls s3://$BUCKET/$FOLDER | awk '{print $4}'`
32+
fi
33+
2234
TEMPFILE=$(mktemp /tmp/index-output.XXXXXX)
2335
PUBLISH_TASK_REPORT=$AIRFLOW_HOME/dags/funcake_dags/scripts/publish_task_report.rb
2436

37+
for record_set in `echo $RESP`
38+
do
39+
if [ -n "$DATA" ]; then
40+
source_key=$record_set
41+
else
42+
source_key=$FOLDER$record_set
43+
fi
44+
bundle exec $INDEXER ingest $(aws s3 presign s3://$BUCKET/$source_key) 2>&1 | tee -a $TEMPFILE
45+
done
2546
report_and_cleanup() {
2647
rc=$?
2748
cat "$TEMPFILE" | ruby "$PUBLISH_TASK_REPORT" || true

tests/index_dag_test.py

Lines changed: 42 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,9 @@ def test_dag_tasks_present(self):
1919
self.assertEqual(self.tasks, [
2020
"harvest_oai",
2121
"create_collection",
22+
"list_index_files",
2223
"combine_index",
24+
"solr_commit",
2325
"solr_alias_swap",
2426
"success",
2527
])
@@ -28,29 +30,45 @@ def test_dag_task_order(self):
2830
"""Unit test that the DAG instance contains the expected dependencies."""
2931
expected_task_deps = {
3032
"create_collection": ["harvest_oai"],
31-
"combine_index": ["create_collection"],
32-
"solr_alias_swap": ["combine_index"],
33+
"list_index_files": ["create_collection"],
34+
"combine_index": ["list_index_files"],
35+
"solr_commit": ["combine_index"],
36+
"solr_alias_swap": ["solr_commit"],
3337
"success": ["solr_alias_swap"],
3438
}
3539

3640
for task, upstream_tasks in expected_task_deps.items():
3741
upstream_list = [up_task.task_id for up_task in FCDAGDEV.get_task(task).upstream_list]
3842
self.assertCountEqual(upstream_tasks, upstream_list)
3943

44+
def test_list_index_files_task(self):
45+
"""Unit test that the DAG instance can list index files from S3."""
46+
task = FCDAGDEV.get_task("list_index_files")
47+
self.assertEqual(task.bucket, "{{ var.value.AIRFLOW_DATA_BUCKET }}")
48+
self.assertEqual(task.prefix, "funcake_dev_index/{{ logical_date.strftime('%Y-%m-%d_%H-%M-%S') }}/new-updated/")
49+
self.assertEqual(task.aws_conn_id, "AIRFLOW_S3")
50+
4051
def test_combine_index_task(self):
4152
"""Unit test that the DAG instance can find required solr indexing bash script."""
4253
task = FCDAGDEV.get_task("combine_index")
4354
expected_bash_path = "{{ var.value.AIRFLOW_HOME }}/dags/funcake_dags/scripts/index.sh "
4455
self.assertEqual(task.bash_command, expected_bash_path)
4556
self.assertEqual(task.env["AIRFLOW_HOME"], "{{ var.value.AIRFLOW_HOME }}")
4657
self.assertEqual(task.env["BUCKET"], "{{ var.value.AIRFLOW_DATA_BUCKET }}")
47-
self.assertEqual(task.env["FOLDER"], "funcake_dev_index/{{ logical_date.strftime('%Y-%m-%d_%H-%M-%S') }}/new-updated/")
58+
self.assertEqual(task.env["DATA"], "{{ ti.xcom_pull(task_ids='list_index_files') | tojson }}")
4859
self.assertEqual(task.env["SOLR_URL"], "{{ conn.get('SOLRCLOUD-WRITER').host if '://' in conn.get('SOLRCLOUD-WRITER').host else 'https://' + conn.get('SOLRCLOUD-WRITER').host }}/solr/{{ var.json.FUNCAKE_SOLR_CONFIG.configset }}-{{ logical_date.strftime('%Y-%m-%d_%H-%M-%S') }}")
4960
self.assertEqual(task.env["SOLR_AUTH_USER"], "{{ conn.get('SOLRCLOUD-WRITER').login or '' }}")
5061
self.assertEqual(task.env["SOLR_AUTH_PASSWORD"], "{{ conn.get('SOLRCLOUD-WRITER').password or '' }}")
5162
self.assertEqual(task.env["AWS_ACCESS_KEY_ID"], "{{ conn.get('AIRFLOW_S3').login }}")
5263
self.assertEqual(task.env["AWS_SECRET_ACCESS_KEY"], "{{ conn.get('AIRFLOW_S3').password }}")
5364

65+
def test_solr_commit_task(self):
66+
"""Unit test that the DAG instance includes a final Solr commit."""
67+
task = FCDAGDEV.get_task("solr_commit")
68+
self.assertEqual(task.http_conn_id, "SOLRCLOUD-WRITER")
69+
self.assertEqual(task.method, "GET")
70+
self.assertEqual(task.endpoint, "/solr/{{ var.json.FUNCAKE_SOLR_CONFIG.configset }}-{{ logical_date.strftime('%Y-%m-%d_%H-%M-%S') }}/update?commit=true")
71+
5472

5573
class TestFuncakeProdIndexDAG(unittest.TestCase):
5674
"""Primary Class for Testing the FunCake Solr Index DAG."""
@@ -68,7 +86,9 @@ def test_dag_tasks_present(self):
6886
self.assertEqual(self.tasks, [
6987
"harvest_oai",
7088
"create_collection",
89+
"list_index_files",
7190
"combine_index",
91+
"solr_commit",
7292
"solr_alias_swap",
7393
"success",
7494
])
@@ -77,25 +97,41 @@ def test_dag_task_order(self):
7797
"""Unit test that the DAG instance contains the expected dependencies."""
7898
expected_task_deps = {
7999
"create_collection": ["harvest_oai"],
80-
"combine_index": ["create_collection"],
81-
"solr_alias_swap": ["combine_index"],
100+
"list_index_files": ["create_collection"],
101+
"combine_index": ["list_index_files"],
102+
"solr_commit": ["combine_index"],
103+
"solr_alias_swap": ["solr_commit"],
82104
"success": ["solr_alias_swap"],
83105
}
84106

85107
for task, upstream_tasks in expected_task_deps.items():
86108
upstream_list = [up_task.task_id for up_task in FCDAGPROD.get_task(task).upstream_list]
87109
self.assertCountEqual(upstream_tasks, upstream_list)
88110

111+
def test_list_index_files_task(self):
112+
"""Unit test that the DAG instance can list index files from S3."""
113+
task = FCDAGPROD.get_task("list_index_files")
114+
self.assertEqual(task.bucket, "{{ var.value.AIRFLOW_DATA_BUCKET }}")
115+
self.assertEqual(task.prefix, "funcake_prod_index/{{ logical_date.strftime('%Y-%m-%d_%H-%M-%S') }}/new-updated/")
116+
self.assertEqual(task.aws_conn_id, "AIRFLOW_S3")
117+
89118
def test_combine_index_task(self):
90119
"""Unit test that the DAG instance can find required solr indexing bash script."""
91120
task = FCDAGPROD.get_task("combine_index")
92121
expected_bash_path = "{{ var.value.AIRFLOW_HOME }}/dags/funcake_dags/scripts/index.sh "
93122
self.assertEqual(task.bash_command, expected_bash_path)
94123
self.assertEqual(task.env["AIRFLOW_HOME"], "{{ var.value.AIRFLOW_HOME }}")
95124
self.assertEqual(task.env["BUCKET"], "{{ var.value.AIRFLOW_DATA_BUCKET }}")
96-
self.assertEqual(task.env["FOLDER"], "funcake_prod_index/{{ logical_date.strftime('%Y-%m-%d_%H-%M-%S') }}/new-updated/")
125+
self.assertEqual(task.env["DATA"], "{{ ti.xcom_pull(task_ids='list_index_files') | tojson }}")
97126
self.assertEqual(task.env["SOLR_URL"], "{{ conn.get('SOLRCLOUD-WRITER').host if '://' in conn.get('SOLRCLOUD-WRITER').host else 'https://' + conn.get('SOLRCLOUD-WRITER').host }}/solr/{{ var.json.FUNCAKE_SOLR_CONFIG.configset }}-{{ logical_date.strftime('%Y-%m-%d_%H-%M-%S') }}")
98127
self.assertEqual(task.env["SOLR_AUTH_USER"], "{{ conn.get('SOLRCLOUD-WRITER').login or '' }}")
99128
self.assertEqual(task.env["SOLR_AUTH_PASSWORD"], "{{ conn.get('SOLRCLOUD-WRITER').password or '' }}")
100129
self.assertEqual(task.env["AWS_ACCESS_KEY_ID"], "{{ conn.get('AIRFLOW_S3').login }}")
101130
self.assertEqual(task.env["AWS_SECRET_ACCESS_KEY"], "{{ conn.get('AIRFLOW_S3').password }}")
131+
132+
def test_solr_commit_task(self):
133+
"""Unit test that the DAG instance includes a final Solr commit."""
134+
task = FCDAGPROD.get_task("solr_commit")
135+
self.assertEqual(task.http_conn_id, "SOLRCLOUD-WRITER")
136+
self.assertEqual(task.method, "GET")
137+
self.assertEqual(task.endpoint, "/solr/{{ var.json.FUNCAKE_SOLR_CONFIG.configset }}-{{ logical_date.strftime('%Y-%m-%d_%H-%M-%S') }}/update?commit=true")

0 commit comments

Comments
 (0)