diff --git a/elementary/monitor/dbt_project/macros/base_queries/current_tests_run_results_query.sql b/elementary/monitor/dbt_project/macros/base_queries/current_tests_run_results_query.sql index 3e18222bf..8c99cd8c1 100644 --- a/elementary/monitor/dbt_project/macros/base_queries/current_tests_run_results_query.sql +++ b/elementary/monitor/dbt_project/macros/base_queries/current_tests_run_results_query.sql @@ -2,14 +2,14 @@ with elementary_test_results as ( select * from {{ ref('elementary', 'elementary_test_results') }} {% if days_back %} - where {{ elementary.edr_datediff(elementary.edr_cast_as_timestamp('detected_at'), elementary.edr_current_timestamp(), 'day') }} < {{ days_back }} + where {{ elementary_cli.days_back_filter('detected_at', days_back, partition_column='created_at') }} {% endif %} ), dbt_run_results as ( select * from {{ ref('elementary', 'dbt_run_results') }} {% if days_back %} - where {{ elementary.edr_datediff(elementary.edr_cast_as_timestamp('execute_completed_at'), elementary.edr_current_timestamp(), 'day') }} < {{ days_back }} + where {{ elementary_cli.days_back_filter('execute_completed_at', days_back, partition_column='created_at', column_is_string=true) }} {% endif %} ), diff --git a/elementary/monitor/dbt_project/macros/get_models_runs.sql b/elementary/monitor/dbt_project/macros/get_models_runs.sql index f17aedd27..a41f4ac7b 100644 --- a/elementary/monitor/dbt_project/macros/get_models_runs.sql +++ b/elementary/monitor/dbt_project/macros/get_models_runs.sql @@ -5,6 +5,7 @@ *, row_number() over (partition by unique_id order by generated_at desc) as invocations_rank_index from {{ ref('elementary', 'model_run_results') }} + where {{ elementary_cli.days_back_filter('generated_at', days_back, partition_column='created_at', column_is_string=true) }} ) select @@ -22,9 +23,8 @@ case when invocations_rank_index = 1 then compiled_code else NULL end as compiled_code, generated_at from model_runs - where {{ elementary.edr_datediff(elementary.edr_cast_as_timestamp('generated_at'), elementary.edr_current_timestamp(), 'day') }} < {{ days_back }} {% if exclude_elementary %} - and unique_id not like 'model.elementary.%' + where unique_id not like 'model.elementary.%' {% endif %} order by generated_at {% endset %} diff --git a/elementary/monitor/dbt_project/macros/get_result_rows_agate.sql b/elementary/monitor/dbt_project/macros/get_result_rows_agate.sql index b56b42126..146bc434b 100644 --- a/elementary/monitor/dbt_project/macros/get_result_rows_agate.sql +++ b/elementary/monitor/dbt_project/macros/get_result_rows_agate.sql @@ -1,14 +1,10 @@ {% macro get_result_rows_agate(days_back, valid_ids_query = none) %} - {% do return(adapter.dispatch('get_result_rows_agate', 'elementary_cli')(days_back, valid_ids_query)) %} -{% endmacro %} - -{% macro default__get_result_rows_agate(days_back, valid_ids_query = none) %} {% set query %} select elementary_test_results_id, result_row from {{ ref("test_result_rows", package="elementary") }} - where {{ elementary.edr_datediff(elementary.edr_cast_as_timestamp('detected_at'), elementary.edr_current_timestamp(), 'day') }} < {{ days_back }} + where {{ elementary_cli.days_back_filter('detected_at', days_back, partition_column='created_at') }} {% if valid_ids_query %} and elementary_test_results_id in ({{ valid_ids_query }}) {% endif %} @@ -19,21 +15,3 @@ {% endif %} {% do return(res.group_by("elementary_test_results_id")) %} {% endmacro %} - -{% macro bigquery__get_result_rows_agate(days_back, valid_ids_query = none) %} - {% set query %} - select - elementary_test_results_id, - result_row - from {{ ref("test_result_rows", package="elementary") }} - where detected_at > {{ elementary.edr_timeadd('day', -1 * days_back, elementary.edr_current_timestamp()) }} - {% if valid_ids_query %} - and elementary_test_results_id in ({{ valid_ids_query }}) - {% endif %} - {% endset %} - {% set res = elementary.run_query(query) %} - {% if not res %} - {% do return({}) %} - {% endif %} - {% do return(res.group_by("elementary_test_results_id")) %} -{% endmacro %} \ No newline at end of file diff --git a/elementary/monitor/dbt_project/macros/utils/days_back_filter.sql b/elementary/monitor/dbt_project/macros/utils/days_back_filter.sql new file mode 100644 index 000000000..c5e544e27 --- /dev/null +++ b/elementary/monitor/dbt_project/macros/utils/days_back_filter.sql @@ -0,0 +1,44 @@ +{# + Filters `column` to the last `days_back` days. + + Pass `partition_column` when `column` is not the table's partition column, or is a string and so + has to be cast: BigQuery prunes only on the raw partition column, and without this the query + scans all history. Its bound carries a day of slack because the two columns are stamped from + different clocks — the guarded ones by the dbt client, `created_at` by the warehouse on insert. +#} + +{% macro days_back_filter(column, days_back, partition_column=none, column_is_string=false) %} + {% do return( + adapter.dispatch("days_back_filter", "elementary_cli")( + column, days_back, partition_column, column_is_string + ) + ) %} +{% endmacro %} + +{% macro default__days_back_filter( + column, days_back, partition_column=none, column_is_string=false +) %} + {%- set days_diff = elementary.edr_datediff( + elementary.edr_cast_as_timestamp(column), elementary.edr_current_timestamp(), "day" + ) -%} + {% do return(days_diff ~ " < " ~ days_back) %} +{% endmacro %} + +{% macro bigquery__days_back_filter( + column, days_back, partition_column=none, column_is_string=false +) %} + {%- set filtered = elementary.edr_cast_as_timestamp(column) if column_is_string else column -%} + {%- set conditions = [ + filtered ~ " > " ~ elementary.edr_timeadd( + "day", -1 * (days_back | int), elementary.edr_current_timestamp() + ) + ] -%} + {%- if partition_column and partition_column != column -%} + {%- do conditions.append( + partition_column ~ " > " ~ elementary.edr_timeadd( + "day", -1 * ((days_back | int) + 1), elementary.edr_current_timestamp() + ) + ) -%} + {%- endif -%} + {% do return(conditions | join(" and ")) %} +{% endmacro %}