From e615bff77f3fdfc46d3fc5770c687dcc7c0e5e5c Mon Sep 17 00:00:00 2001 From: Arturo Volpe Date: Mon, 21 Dec 2020 21:37:54 -0300 Subject: [PATCH 1/5] feat: Add ETL to re-index the elasticsearch indices Signed-off-by: Arturo Volpe --- infra/elk/logstash_jdbc_input/Dockerfile | 2 +- .../fts_authorities_ddjj.conf | 33 ++++++++++ .../logstash_jdbc_input/fts_full_data.conf | 24 +++++++ infra/elk/logstash_jdbc_input/logstash.conf | 64 ------------------- infra/elk/re-index.sh | 24 +++++++ .../airflow/dags/elastic_fts_auth_ddjj.py | 42 ++++++++++++ .../dags/elastic_fts_full_data_index.py | 43 +++++++++++++ .../dags/sql/elastic_index_fts_auth_ddjj.sql | 50 +++++++++++++++ .../dags/sql/elastic_index_full_data.sql | 2 +- 9 files changed, 218 insertions(+), 66 deletions(-) create mode 100644 infra/elk/logstash_jdbc_input/fts_authorities_ddjj.conf create mode 100644 infra/elk/logstash_jdbc_input/fts_full_data.conf delete mode 100644 infra/elk/logstash_jdbc_input/logstash.conf create mode 100644 infra/elk/re-index.sh create mode 100644 scripts/python/airflow/dags/elastic_fts_auth_ddjj.py create mode 100644 scripts/python/airflow/dags/elastic_fts_full_data_index.py create mode 100644 scripts/python/airflow/dags/sql/elastic_index_fts_auth_ddjj.sql rename migrations/R_Full_data_table.sql => scripts/python/airflow/dags/sql/elastic_index_full_data.sql (98%) diff --git a/infra/elk/logstash_jdbc_input/Dockerfile b/infra/elk/logstash_jdbc_input/Dockerfile index 012226f..5ce0a17 100644 --- a/infra/elk/logstash_jdbc_input/Dockerfile +++ b/infra/elk/logstash_jdbc_input/Dockerfile @@ -8,5 +8,5 @@ RUN /usr/share/logstash/bin/logstash-plugin install logstash-filter-aggregate && RUN curl https://jdbc.postgresql.org/download/postgresql-42.2.18.jar -o $LOGSTASH_JDBC_DRIVER_JAR_LOCATION -ADD ./logstash.conf /usr/share/logstash/pipeline/logstash.conf +COPY ./*.conf /usr/share/logstash/pipeline/ diff --git a/infra/elk/logstash_jdbc_input/fts_authorities_ddjj.conf b/infra/elk/logstash_jdbc_input/fts_authorities_ddjj.conf new file mode 100644 index 0000000..30636c4 --- /dev/null +++ b/infra/elk/logstash_jdbc_input/fts_authorities_ddjj.conf @@ -0,0 +1,33 @@ +input { + jdbc { + # AQE + jdbc_driver_library => "${LOGSTASH_JDBC_DRIVER_JAR_LOCATION}" + jdbc_driver_class => "${LOGSTASH_JDBC_DRIVER}" + jdbc_connection_string => "${LOGSTASH_JDBC_URL}" + jdbc_user => "${LOGSTASH_JDBC_USERNAME}" + jdbc_password => "${LOGSTASH_JDBC_PASSWORD}" + schedule => "* * * * *" + use_column_value => true + tracking_column => "id" + statement => 'SELECT id, document, last_name, first_name, year_elected, departament, charge, list, title, sex, nacionality, age, CAST(start as TEXT), CAST("end" as TEXT), presented FROM analysis.authorities_with_ddjj WHERE id > :sql_last_value ORDER BY id' + last_run_metadata_path => "/logstash_data/ddjj.yml" + } +} +filter { + json { + source => "start" + target => "start" + } + json { + source => "end" + target => "end" + } +} +output { + elasticsearch { + hosts => ["${LOGSTASH_ELASTICSEARCH_HOST}"] + index => "fts_authorities_ddjj" + document_id => "%{id}" + } + stdout { codec => json_lines } +} diff --git a/infra/elk/logstash_jdbc_input/fts_full_data.conf b/infra/elk/logstash_jdbc_input/fts_full_data.conf new file mode 100644 index 0000000..2dec294 --- /dev/null +++ b/infra/elk/logstash_jdbc_input/fts_full_data.conf @@ -0,0 +1,24 @@ +input { + jdbc { + # full_data + jdbc_driver_library => "${LOGSTASH_JDBC_DRIVER_JAR_LOCATION}" + jdbc_driver_class => "${LOGSTASH_JDBC_DRIVER}" + jdbc_connection_string => "${LOGSTASH_JDBC_URL}" + jdbc_user => "${LOGSTASH_JDBC_USERNAME}" + jdbc_password => "${LOGSTASH_JDBC_PASSWORD}" + schedule => "* * * * *" + use_column_value => true + tracking_column => "document" + statement => "select * from analysis.full_data WHERE document > :sql_last_value ORDER BY document" + last_run_metadata_path => "/logstash_data/full_data-jdbc-int-sql_last_value.yml" + } +} + +output { + elasticsearch { + hosts => ["${LOGSTASH_ELASTICSEARCH_HOST}"] + index => "fts_full_data" + document_id => "%{document}" + } + stdout { codec => json_lines } +} diff --git a/infra/elk/logstash_jdbc_input/logstash.conf b/infra/elk/logstash_jdbc_input/logstash.conf deleted file mode 100644 index 3d66976..0000000 --- a/infra/elk/logstash_jdbc_input/logstash.conf +++ /dev/null @@ -1,64 +0,0 @@ -input { - jdbc { - # AQE - jdbc_driver_library => "${LOGSTASH_JDBC_DRIVER_JAR_LOCATION}" - jdbc_driver_class => "${LOGSTASH_JDBC_DRIVER}" - jdbc_connection_string => "${LOGSTASH_JDBC_URL}" - jdbc_user => "${LOGSTASH_JDBC_USERNAME}" - jdbc_password => "${LOGSTASH_JDBC_PASSWORD}" - schedule => "* 0 * * *" - statement => "SELECT regexp_replace(identifier, '\.', '', 'g') as document, name || ' ' || lastname as name, 'a_quien_elegimos' as source, head_shot as photo FROM staging.a_quien_elegimos" - last_run_metadata_path => "/logstash_data/aqe-jdbc-int-sql_last_value.yml" - } - - jdbc { - # Declarations - jdbc_driver_library => "${LOGSTASH_JDBC_DRIVER_JAR_LOCATION}" - jdbc_driver_class => "${LOGSTASH_JDBC_DRIVER}" - jdbc_connection_string => "${LOGSTASH_JDBC_URL}" - jdbc_user => "${LOGSTASH_JDBC_USERNAME}" - jdbc_password => "${LOGSTASH_JDBC_PASSWORD}" - schedule => "* * * * *" - use_column_value => true - tracking_column => "id" - statement => "SELECT id, 'declarations' as source, document, name, net_worth, active, passive FROM analysis.declarations WHERE id > :sql_last_value" - last_run_metadata_path => "/logstash_data/declarations-jdbc-int-sql_last_value.yml" - } - - jdbc { - # ANDE - jdbc_driver_library => "${LOGSTASH_JDBC_DRIVER_JAR_LOCATION}" - jdbc_driver_class => "${LOGSTASH_JDBC_DRIVER}" - jdbc_connection_string => "${LOGSTASH_JDBC_URL}" - jdbc_user => "${LOGSTASH_JDBC_USERNAME}" - jdbc_password => "${LOGSTASH_JDBC_PASSWORD}" - schedule => "* * * * *" - use_column_value => true - tracking_column => "id" - statement => "SELECT id, documento as document, cliente as name, 'ande_exonerados' as source FROM staging.ande_exonerados WHERE id > :sql_last_value" - last_run_metadata_path => "/logstash_data/ande_exonerados-jdbc-int-sql_last_value.yml" - } - - jdbc { - # TSJE - jdbc_driver_library => "${LOGSTASH_JDBC_DRIVER_JAR_LOCATION}" - jdbc_driver_class => "${LOGSTASH_JDBC_DRIVER}" - jdbc_connection_string => "${LOGSTASH_JDBC_URL}" - jdbc_user => "${LOGSTASH_JDBC_USERNAME}" - jdbc_password => "${LOGSTASH_JDBC_PASSWORD}" - schedule => "* * * * *" - use_column_value => true - tracking_column => "id" - statement => "SELECT id, nombre || ' ' || apellido as name, cedula as document, edad as age, 'tsje_elected' as source FROM analysis.tsje_elected WHERE id > :sql_last_value" - last_run_metadata_path => "/logstash_data/tsje_elected-jdbc-int-sql_last_value.yml" - } -} - -output { - elasticsearch { - hosts => ["${LOGSTASH_ELASTICSEARCH_HOST}"] - index => "fts_people" - document_id => "%{id}" - } - stdout { codec => json_lines } -} diff --git a/infra/elk/re-index.sh b/infra/elk/re-index.sh new file mode 100644 index 0000000..a1dbbfa --- /dev/null +++ b/infra/elk/re-index.sh @@ -0,0 +1,24 @@ +#!/bin/bash + + +set -e +set -x + +INDEX_NAME=$1 +FILE_NAME=$2 + + +docker-compose stop logstash + +rm -rvf "logstash_data/$FILE_NAME" + +docker-compose exec elasticsearch curl -XDELETE localhost:9200/$INDEX_NAME + +echo "Index removed, restarting" + + +docker-compose up -d --build logstash + +echo "The index should start updating very soon, check logs" + + diff --git a/scripts/python/airflow/dags/elastic_fts_auth_ddjj.py b/scripts/python/airflow/dags/elastic_fts_auth_ddjj.py new file mode 100644 index 0000000..0239f61 --- /dev/null +++ b/scripts/python/airflow/dags/elastic_fts_auth_ddjj.py @@ -0,0 +1,42 @@ +from datetime import timedelta, datetime + +from airflow import DAG +from airflow.operators.bash_operator import BashOperator +from airflow.operators.postgres_operator import PostgresOperator + +default_args = { + 'owner': 'airflow', + 'depends_on_past': False, + 'email': ['arturovolpe@gmail.com'], + 'email_on_failure': False, + 'email_on_retry': False, + 'retries': 1, + 'retry_delay': timedelta(seconds=5), + 'params': { + } +} +dag = DAG( + 'elastic_fts_auth_ddjj', + default_args=default_args, + description='ETL that creates the table for search authorities with ddjj and then indexes the table in a elastic', + start_date=datetime(2020, 12, 21), + schedule_interval=timedelta(weeks=1), +) + +with dag: + do_curl = BashOperator( + task_id=f'call_webhook', + bash_command=f""" + curl {{ var.value.ELASTIC_IDX_AUTH_DDJJ }} + """, + retries=10 + ) + + do_query = PostgresOperator(task_id='do_query', + sql="sql/elastic_index_fts_auth_ddjj.sql") + + do_query >> do_curl + +if __name__ == '__main__': + dag.clear(reset_dag_runs=True) + dag.run() diff --git a/scripts/python/airflow/dags/elastic_fts_full_data_index.py b/scripts/python/airflow/dags/elastic_fts_full_data_index.py new file mode 100644 index 0000000..94e8c43 --- /dev/null +++ b/scripts/python/airflow/dags/elastic_fts_full_data_index.py @@ -0,0 +1,43 @@ +from datetime import timedelta, datetime + +from airflow import DAG +from airflow.operators.bash_operator import BashOperator +from airflow.operators.postgres_operator import PostgresOperator + + +default_args = { + 'owner': 'airflow', + 'depends_on_past': False, + 'email': ['arturovolpe@gmail.com'], + 'email_on_failure': False, + 'email_on_retry': False, + 'retries': 1, + 'retry_delay': timedelta(seconds=5), + 'params': { + } +} +dag = DAG( + 'elastic_fts_full_data_index', + default_args=default_args, + description='ETL that creates the table with all people data', + start_date=datetime(2020, 12, 21), + schedule_interval=timedelta(weeks=1), +) + +with dag: + do_curl = BashOperator( + task_id=f'call_webhook', + bash_command=f""" + curl {{ var.value.ELASTIC_IDX_FULL_DATA_HOOK }} + """, + retries=10 + ) + + do_query = PostgresOperator(task_id='do_query', + sql="sql/elastic_index_full_data.sql") + + do_query >> do_curl + +if __name__ == '__main__': + dag.clear(reset_dag_runs=True) + dag.run() diff --git a/scripts/python/airflow/dags/sql/elastic_index_fts_auth_ddjj.sql b/scripts/python/airflow/dags/sql/elastic_index_fts_auth_ddjj.sql new file mode 100644 index 0000000..985445a --- /dev/null +++ b/scripts/python/airflow/dags/sql/elastic_index_fts_auth_ddjj.sql @@ -0,0 +1,50 @@ +DROP TABLE analysis.authorities_with_ddjj; + +CREATE TABLE analysis.authorities_with_ddjj AS ( + SELECT autoridad.id, + autoridad.cedula AS document, + autoridad.apellido AS last_name, + autoridad.nombre AS first_name, + autoridad.ano AS year_elected, + autoridad.dep_desc AS departament, + autoridad.cand_desc AS charge, + autoridad.nombre_lista AS list, + autoridad.desc_tit_sup AS title, + autoridad.sexo AS sex, + autoridad.nacionalidad AS nacionality, + autoridad.edad AS age, + (SELECT jsonb_build_object( + 'id', start.id, + 'link', start.link, + 'link_sandwich', start.link_sandwich, + 'origin', start.origin, + 'active', start.active, + 'passive', start.passive, + 'net_worth', start.net_worth) + FROM analysis.declarations start + WHERE start.document = autoridad.cedula + AND start.year = autoridad.ano + ORDER BY start.version desc + LIMIT 1) AS start, + (SELECT jsonb_build_object( + 'id', "end".id, + 'link', "end".link, + 'link_sandwich', "end".link_sandwich, + 'origin', "end".origin, + 'active', "end".active, + 'passive', "end".passive, + 'net_worth', "end".net_worth) + FROM analysis.declarations "end" + WHERE "end".document = autoridad.cedula + AND "end".year = autoridad.ano + 5 + ORDER BY "end".version desc + LIMIT 1) AS "end" + FROM analysis.tsje_elected autoridad +); + +ALTER TABLE analysis.authorities_with_ddjj + ADD COLUMN presented boolean DEFAULT FALSE; + +UPDATE analysis.authorities_with_ddjj +SET presented = start IS NOT NULL AND "end" IS NOT NULL +; diff --git a/migrations/R_Full_data_table.sql b/scripts/python/airflow/dags/sql/elastic_index_full_data.sql similarity index 98% rename from migrations/R_Full_data_table.sql rename to scripts/python/airflow/dags/sql/elastic_index_full_data.sql index 4f585fa..da9b97a 100644 --- a/migrations/R_Full_data_table.sql +++ b/scripts/python/airflow/dags/sql/elastic_index_full_data.sql @@ -123,7 +123,7 @@ CREATE TABLE analysis.full_data AS ( (raw.ano || to_char(raw.mes, '00')) = docs.last_month_worked GROUP BY docs.document, docs.last_month_worked ) - SELECT CAST(regexp_replace(document, '[^0-9]+', '', 'g') as bigint) as document, + SELECT regexp_replace(document, '[^0-9]+', '', 'g') as document, array_agg(name) as name, array_agg(photo) as photo, array_agg(salary) as salary, From 13dc381c4ae0a7cfed37bcf2268d68fb436864f6 Mon Sep 17 00:00:00 2001 From: Arturo Volpe Date: Mon, 21 Dec 2020 21:46:45 -0300 Subject: [PATCH 2/5] feat: Add photo to index of authorities_with_ddjj Signed-off-by: Arturo Volpe --- infra/elk/logstash_jdbc_input/fts_authorities_ddjj.conf | 2 +- .../python/airflow/dags/sql/elastic_index_fts_auth_ddjj.sql | 4 +++- 2 files changed, 4 insertions(+), 2 deletions(-) diff --git a/infra/elk/logstash_jdbc_input/fts_authorities_ddjj.conf b/infra/elk/logstash_jdbc_input/fts_authorities_ddjj.conf index 30636c4..cde22f2 100644 --- a/infra/elk/logstash_jdbc_input/fts_authorities_ddjj.conf +++ b/infra/elk/logstash_jdbc_input/fts_authorities_ddjj.conf @@ -9,7 +9,7 @@ input { schedule => "* * * * *" use_column_value => true tracking_column => "id" - statement => 'SELECT id, document, last_name, first_name, year_elected, departament, charge, list, title, sex, nacionality, age, CAST(start as TEXT), CAST("end" as TEXT), presented FROM analysis.authorities_with_ddjj WHERE id > :sql_last_value ORDER BY id' + statement => 'SELECT id, document, last_name, first_name, year_elected, departament, charge, list, title, sex, nacionality, age, CAST(start as TEXT), CAST("end" as TEXT), presented, photo FROM analysis.authorities_with_ddjj WHERE id > :sql_last_value ORDER BY id' last_run_metadata_path => "/logstash_data/ddjj.yml" } } diff --git a/scripts/python/airflow/dags/sql/elastic_index_fts_auth_ddjj.sql b/scripts/python/airflow/dags/sql/elastic_index_fts_auth_ddjj.sql index 985445a..a5010d9 100644 --- a/scripts/python/airflow/dags/sql/elastic_index_fts_auth_ddjj.sql +++ b/scripts/python/airflow/dags/sql/elastic_index_fts_auth_ddjj.sql @@ -38,8 +38,10 @@ CREATE TABLE analysis.authorities_with_ddjj AS ( WHERE "end".document = autoridad.cedula AND "end".year = autoridad.ano + 5 ORDER BY "end".version desc - LIMIT 1) AS "end" + LIMIT 1) AS "end", + aqe.head_shot as photo FROM analysis.tsje_elected autoridad + LEFT JOIN staging.a_quien_elegimos aqe ON aqe.identifier = autoridad.cedula ); ALTER TABLE analysis.authorities_with_ddjj From 4ee86bb9547c6b171371f409008286c4e31f3eda Mon Sep 17 00:00:00 2001 From: Arturo Volpe Date: Mon, 21 Dec 2020 21:50:13 -0300 Subject: [PATCH 3/5] feat: Allow re-index to be executed as a webhook Signed-off-by: Arturo Volpe --- infra/elk/re-index.sh | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/infra/elk/re-index.sh b/infra/elk/re-index.sh index a1dbbfa..c204673 100644 --- a/infra/elk/re-index.sh +++ b/infra/elk/re-index.sh @@ -12,7 +12,7 @@ docker-compose stop logstash rm -rvf "logstash_data/$FILE_NAME" -docker-compose exec elasticsearch curl -XDELETE localhost:9200/$INDEX_NAME +docker-compose exec -T elasticsearch curl -XDELETE localhost:9200/$INDEX_NAME echo "Index removed, restarting" From ede4e450248f27f23b275c1aa25c793047373161 Mon Sep 17 00:00:00 2001 From: Arturo Volpe Date: Mon, 21 Dec 2020 22:29:47 -0300 Subject: [PATCH 4/5] fix: Call to webhook on elastic dags Signed-off-by: Arturo Volpe --- scripts/python/airflow/dags/elastic_fts_auth_ddjj.py | 4 ++-- .../airflow/dags/elastic_fts_full_data_index.py | 11 +++++++---- .../airflow/dags/sql/elastic_index_full_data.sql | 1 - 3 files changed, 9 insertions(+), 7 deletions(-) diff --git a/scripts/python/airflow/dags/elastic_fts_auth_ddjj.py b/scripts/python/airflow/dags/elastic_fts_auth_ddjj.py index 0239f61..f38bdf7 100644 --- a/scripts/python/airflow/dags/elastic_fts_auth_ddjj.py +++ b/scripts/python/airflow/dags/elastic_fts_auth_ddjj.py @@ -26,8 +26,8 @@ with dag: do_curl = BashOperator( task_id=f'call_webhook', - bash_command=f""" - curl {{ var.value.ELASTIC_IDX_AUTH_DDJJ }} + bash_command=""" + curl "{{ var.value.ELASTIC_IDX_AUTH_DDJJ }}" """, retries=10 ) diff --git a/scripts/python/airflow/dags/elastic_fts_full_data_index.py b/scripts/python/airflow/dags/elastic_fts_full_data_index.py index 94e8c43..c93da3d 100644 --- a/scripts/python/airflow/dags/elastic_fts_full_data_index.py +++ b/scripts/python/airflow/dags/elastic_fts_full_data_index.py @@ -4,7 +4,6 @@ from airflow.operators.bash_operator import BashOperator from airflow.operators.postgres_operator import PostgresOperator - default_args = { 'owner': 'airflow', 'depends_on_past': False, @@ -27,16 +26,20 @@ with dag: do_curl = BashOperator( task_id=f'call_webhook', - bash_command=f""" - curl {{ var.value.ELASTIC_IDX_FULL_DATA_HOOK }} + bash_command=""" + curl "{{ var.value.ELASTIC_IDX_FULL_DATA_HOOK }}" """, retries=10 ) + clean_db = PostgresOperator(task_id='clean_table', + sql="DROP TABLE analysis.full_data") + do_query = PostgresOperator(task_id='do_query', + sql="sql/elastic_index_full_data.sql") - do_query >> do_curl + clean_db >> do_query >> do_curl if __name__ == '__main__': dag.clear(reset_dag_runs=True) diff --git a/scripts/python/airflow/dags/sql/elastic_index_full_data.sql b/scripts/python/airflow/dags/sql/elastic_index_full_data.sql index da9b97a..e9efb87 100644 --- a/scripts/python/airflow/dags/sql/elastic_index_full_data.sql +++ b/scripts/python/airflow/dags/sql/elastic_index_full_data.sql @@ -1,4 +1,3 @@ -DROP TABLE analysis.full_data; CREATE TABLE analysis.full_data AS ( WITH sfp_documents AS ( select documento as document, From 2488cdb09ec71b50a14e37d858fc05b433670843 Mon Sep 17 00:00:00 2001 From: Rodrigo Benitez Date: Tue, 22 Dec 2020 10:32:35 -0300 Subject: [PATCH 5/5] fix: Change condition for presented declarations --- scripts/python/airflow/dags/sql/elastic_index_fts_auth_ddjj.sql | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/scripts/python/airflow/dags/sql/elastic_index_fts_auth_ddjj.sql b/scripts/python/airflow/dags/sql/elastic_index_fts_auth_ddjj.sql index a5010d9..0a9cb04 100644 --- a/scripts/python/airflow/dags/sql/elastic_index_fts_auth_ddjj.sql +++ b/scripts/python/airflow/dags/sql/elastic_index_fts_auth_ddjj.sql @@ -48,5 +48,5 @@ ALTER TABLE analysis.authorities_with_ddjj ADD COLUMN presented boolean DEFAULT FALSE; UPDATE analysis.authorities_with_ddjj -SET presented = start IS NOT NULL AND "end" IS NOT NULL +SET presented = start IS NOT NULL AND ("end" IS NOT NULL OR year + 5 > date_part('year', CURRENT_DATE)); ;