From 1843290b338bd3bc58388bc5c53ea1611668759d Mon Sep 17 00:00:00 2001 From: bobbyiliev Date: Thu, 6 Aug 2026 11:48:41 +0300 Subject: [PATCH] dbt-materialize: remove duplicated readiness SQL and deploy boilerplate --- .../dbt/adapters/materialize/impl.py | 27 +--- .../macros/deploy/deploy_await.sql | 9 +- .../macros/deploy/deploy_cleanup.sql | 8 +- .../macros/deploy/deploy_helpers.sql | 29 ++++ .../materialize/macros/deploy/deploy_init.sql | 8 +- .../macros/deploy/deploy_promote.sql | 8 +- .../macros/materializations/seed/helpers.sql | 2 +- .../macros/utils/is_cluster_ready.sql | 147 +----------------- 8 files changed, 41 insertions(+), 197 deletions(-) create mode 100644 misc/dbt-materialize/dbt/include/materialize/macros/deploy/deploy_helpers.sql diff --git a/misc/dbt-materialize/dbt/adapters/materialize/impl.py b/misc/dbt-materialize/dbt/adapters/materialize/impl.py index f42491632526e..9191375fb7d61 100644 --- a/misc/dbt-materialize/dbt/adapters/materialize/impl.py +++ b/misc/dbt-materialize/dbt/adapters/materialize/impl.py @@ -15,9 +15,8 @@ # limitations under the License. import subprocess import time -from collections import namedtuple from dataclasses import dataclass -from typing import Any, Dict, List, Optional, Set +from typing import Any, Dict, List, Optional import dbt_common.exceptions import psycopg2 @@ -156,30 +155,6 @@ def _link_cached_relations(self, manifest): # [0]: https://github.com/dbt-labs/dbt-core/blob/13b18654f03d92eab3f5a9113e526a2a844f145d/plugins/postgres/dbt/adapters/postgres/impl.py#L126-L133 pass - def _link_cached_database_relations(self, schemas: Set[str]): - """ - :param schemas: The set of schemas that should have links added. - """ - database = self.config.credentials.database - _Relation = namedtuple("_Relation", "database schema identifier") - links = [ - ( - _Relation(database, dep_schema, dep_identifier), - _Relation(database, ref_schema, ref_identifier), - ) - for dep_schema, dep_identifier, ref_schema, ref_identifier in self.execute_macro( - "materialize__get_relations" - ) - # don't record in cache if this relation isn't in a relevant schema - if ref_schema in schemas - ] - - for dependent, referenced in links: - self.cache.add_link( - referenced=self.Relation.create(**referenced._asdict()), - dependent=self.Relation.create(**dependent._asdict()), - ) - def verify_database(self, database): pass diff --git a/misc/dbt-materialize/dbt/include/materialize/macros/deploy/deploy_await.sql b/misc/dbt-materialize/dbt/include/materialize/macros/deploy/deploy_await.sql index 7418f45f2801f..7c33926feff72 100644 --- a/misc/dbt-materialize/dbt/include/materialize/macros/deploy/deploy_await.sql +++ b/misc/dbt-materialize/dbt/include/materialize/macros/deploy/deploy_await.sql @@ -33,14 +33,7 @@ until all deployment clusters' objects are fully hydrated. #} -{% set current_target_name = target.name %} -{% set deployment = var('deployment') %} -{% set target_config = deployment[current_target_name] %} - --- Check if the target-specific configuration exists -{% if not target_config %} - {{ exceptions.raise_compiler_error("No deployment configuration found for target " ~ current_target_name) }} -{% endif %} +{% set target_config = internal_get_deployment_config() %} {% set clusters = target_config.get('clusters', []) %} diff --git a/misc/dbt-materialize/dbt/include/materialize/macros/deploy/deploy_cleanup.sql b/misc/dbt-materialize/dbt/include/materialize/macros/deploy/deploy_cleanup.sql index 0f2711890905e..df879040b883e 100644 --- a/misc/dbt-materialize/dbt/include/materialize/macros/deploy/deploy_cleanup.sql +++ b/misc/dbt-materialize/dbt/include/materialize/macros/deploy/deploy_cleanup.sql @@ -16,13 +16,7 @@ {% macro deploy_cleanup() %} {% set current_target_name = target.name %} -{% set deployment = var('deployment') %} -{% set target_config = deployment[current_target_name] %} - --- Check if the target-specific configuration exists -{% if not target_config %} - {{ exceptions.raise_compiler_error("No deployment configuration found for target " ~ current_target_name) }} -{% endif %} +{% set target_config = internal_get_deployment_config() %} {{ log("Dropping deployment environment for target " ~ current_target_name, info=True) }} diff --git a/misc/dbt-materialize/dbt/include/materialize/macros/deploy/deploy_helpers.sql b/misc/dbt-materialize/dbt/include/materialize/macros/deploy/deploy_helpers.sql new file mode 100644 index 0000000000000..790bfba5b7673 --- /dev/null +++ b/misc/dbt-materialize/dbt/include/materialize/macros/deploy/deploy_helpers.sql @@ -0,0 +1,29 @@ +-- Copyright Materialize, Inc. and contributors. All rights reserved. +-- +-- Licensed under the Apache License, Version 2.0 (the "License"); +-- you may not use this file except in compliance with the License. +-- You may obtain a copy of the License in the LICENSE file at the +-- root of this repository, or online at +-- +-- http://www.apache.org/licenses/LICENSE-2.0 +-- +-- Unless required by applicable law or agreed to in writing, software +-- distributed under the License is distributed on an "AS IS" BASIS, +-- WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +-- See the License for the specific language governing permissions and +-- limitations under the License. + +{# +Returns the deployment configuration for the current target, raising a compiler +error when the `deployment` variable has no entry for it. Every deploy operation +starts from this configuration. +#} +{% macro internal_get_deployment_config() %} + {% set target_config = var('deployment')[target.name] %} + + {% if not target_config %} + {{ exceptions.raise_compiler_error("No deployment configuration found for target " ~ target.name) }} + {% endif %} + + {{ return(target_config) }} +{% endmacro %} diff --git a/misc/dbt-materialize/dbt/include/materialize/macros/deploy/deploy_init.sql b/misc/dbt-materialize/dbt/include/materialize/macros/deploy/deploy_init.sql index acea343ff5ac4..7959352b831c1 100644 --- a/misc/dbt-materialize/dbt/include/materialize/macros/deploy/deploy_init.sql +++ b/misc/dbt-materialize/dbt/include/materialize/macros/deploy/deploy_init.sql @@ -16,13 +16,7 @@ {% macro deploy_init(ignore_existing_objects=False) %} {% set current_target_name = target.name %} -{% set deployment = var('deployment') %} -{% set target_config = deployment[current_target_name] %} - --- Check if the target-specific configuration exists -{% if not target_config %} - {{ exceptions.raise_compiler_error("No deployment configuration found for target " ~ current_target_name) }} -{% endif %} +{% set target_config = internal_get_deployment_config() %} {{ log("Creating deployment environment for target " ~ current_target_name, info=True) }} diff --git a/misc/dbt-materialize/dbt/include/materialize/macros/deploy/deploy_promote.sql b/misc/dbt-materialize/dbt/include/materialize/macros/deploy/deploy_promote.sql index 424da24ab649e..61df11a4fadda 100644 --- a/misc/dbt-materialize/dbt/include/materialize/macros/deploy/deploy_promote.sql +++ b/misc/dbt-materialize/dbt/include/materialize/macros/deploy/deploy_promote.sql @@ -49,13 +49,7 @@ #} {% set current_target_name = target.name %} -{% set deployment = var('deployment') %} -{% set target_config = deployment[current_target_name] %} - --- Check if the target-specific configuration exists -{% if not target_config %} - {{ exceptions.raise_compiler_error("No deployment configuration found for target " ~ current_target_name) }} -{% endif %} +{% set target_config = internal_get_deployment_config() %} {{ log("Creating deployment environment for target " ~ current_target_name, info=True) }} diff --git a/misc/dbt-materialize/dbt/include/materialize/macros/materializations/seed/helpers.sql b/misc/dbt-materialize/dbt/include/materialize/macros/materializations/seed/helpers.sql index 45fa0ba1bc5c8..04d7fc87f1f82 100644 --- a/misc/dbt-materialize/dbt/include/materialize/macros/materializations/seed/helpers.sql +++ b/misc/dbt-materialize/dbt/include/materialize/macros/materializations/seed/helpers.sql @@ -13,7 +13,7 @@ -- See the License for the specific language governing permissions and -- limitations under the License. -{% macro materialize__reset_csv_table(model, full_refresh, old_relation, agate_table, cluster) %} +{% macro materialize__reset_csv_table(model, full_refresh, old_relation, agate_table) %} {% set sql = "" %} -- Allow setting a cluster configuration for seeds in `dbt_project.yml` or -- a .yml file in the seed target path. If no cluster is configured, use diff --git a/misc/dbt-materialize/dbt/include/materialize/macros/utils/is_cluster_ready.sql b/misc/dbt-materialize/dbt/include/materialize/macros/utils/is_cluster_ready.sql index 70975ab7fcf91..358c906cf325e 100644 --- a/misc/dbt-materialize/dbt/include/materialize/macros/utils/is_cluster_ready.sql +++ b/misc/dbt-materialize/dbt/include/materialize/macros/utils/is_cluster_ready.sql @@ -39,147 +39,12 @@ {{ exceptions.raise_compiler_error("No cluster specified and no default cluster found in target profile.") }} {% endif %} -{{ log("Checking cluster readiness for: " ~ cluster, info=True) }} +{% set statuses = are_clusters_ready([cluster], lag_threshold) %} -{%- set check_cluster_ready_sql %} -WITH --- Detect problematic replicas: 3+ OOM kills in 24h -problematic_replicas AS ( - SELECT replica_id - FROM mz_internal.mz_cluster_replica_status_history - WHERE occurred_at + INTERVAL '24 hours' > mz_now() - AND reason = 'oom-killed' - GROUP BY replica_id - HAVING COUNT(*) >= 3 -), - --- Cluster health: count total vs problematic replicas -cluster_health AS ( - SELECT - c.name AS cluster_name, - c.id AS cluster_id, - COUNT(r.id) AS total_replicas, - COUNT(pr.replica_id) AS problematic_replicas - FROM mz_clusters c - LEFT JOIN mz_cluster_replicas r ON c.id = r.cluster_id - LEFT JOIN problematic_replicas pr ON r.id = pr.replica_id - WHERE c.name = {{ dbt.string_literal(cluster) }} - GROUP BY c.name, c.id -), - --- Hydration counts per cluster (best replica) -hydration_counts AS ( - SELECT - c.name AS cluster_name, - r.id AS replica_id, - COUNT(*) FILTER (WHERE mhs.hydrated) AS hydrated, - COUNT(*) AS total - FROM mz_clusters c - JOIN mz_cluster_replicas r ON c.id = r.cluster_id - LEFT JOIN mz_internal.mz_hydration_statuses mhs ON mhs.replica_id = r.id - WHERE c.name = {{ dbt.string_literal(cluster) }} - GROUP BY c.name, r.id -), - -hydration_best AS ( - SELECT cluster_name, MAX(hydrated) AS hydrated, MAX(total) AS total - FROM hydration_counts - GROUP BY cluster_name -), - --- Max lag per cluster using mz_wallclock_global_lag -cluster_lag AS ( - SELECT - c.name AS cluster_name, - MAX(EXTRACT(EPOCH FROM wgl.lag)) AS max_lag_secs - FROM mz_clusters c - JOIN mz_cluster_replicas r ON c.id = r.cluster_id - JOIN mz_internal.mz_hydration_statuses mhs ON mhs.replica_id = r.id - JOIN mz_internal.mz_wallclock_global_lag wgl ON wgl.object_id = mhs.object_id - WHERE c.name = {{ dbt.string_literal(cluster) }} - GROUP BY c.name -), - --- Convert lag_threshold interval to seconds -lag_threshold_secs AS ( - SELECT EXTRACT(EPOCH FROM INTERVAL {{ dbt.string_literal(lag_threshold) }}) AS threshold_secs -) - -SELECT - ch.cluster_name, - ch.cluster_id, - CASE - WHEN ch.total_replicas = 0 THEN 'failing' - WHEN ch.total_replicas = ch.problematic_replicas THEN 'failing' - WHEN COALESCE(hb.hydrated, 0) < COALESCE(hb.total, 0) THEN 'hydrating' - WHEN COALESCE(cl.max_lag_secs, 0) > (SELECT threshold_secs FROM lag_threshold_secs) THEN 'lagging' - ELSE 'ready' - END AS status, - CASE - WHEN ch.total_replicas = 0 THEN 'no_replicas' - WHEN ch.total_replicas = ch.problematic_replicas THEN 'all_replicas_problematic' - ELSE NULL - END AS failure_reason, - COALESCE(hb.hydrated, 0) AS hydrated_count, - COALESCE(hb.total, 0) AS total_count, - COALESCE(cl.max_lag_secs, 0)::bigint AS max_lag_secs, - ch.total_replicas, - ch.problematic_replicas, - (SELECT threshold_secs FROM lag_threshold_secs)::bigint AS threshold_secs -FROM cluster_health ch -LEFT JOIN hydration_best hb ON ch.cluster_name = hb.cluster_name -LEFT JOIN cluster_lag cl ON ch.cluster_name = cl.cluster_name -{%- endset %} - -{%- set results = run_query(check_cluster_ready_sql) %} -{%- if execute -%} - {%- if results and results.rows|length > 0 -%} - {%- set row = results.rows[0] -%} - {%- set cluster_name = row[0] -%} - {%- set cluster_id = row[1] -%} - {%- set status = row[2] -%} - {%- set failure_reason = row[3] -%} - {%- set hydrated_count = row[4] -%} - {%- set total_count = row[5] -%} - {%- set max_lag_secs = row[6] -%} - {%- set total_replicas = row[7] -%} - {%- set problematic_replicas = row[8] -%} - {%- set threshold_secs = row[9] -%} - - {#- Log status details -#} - {%- if status == 'ready' -%} - {{ log("Cluster " ~ cluster ~ " is ready. Hydration: " ~ hydrated_count ~ "/" ~ total_count ~ ", Lag: " ~ max_lag_secs ~ "s", info=True) }} - {%- elif status == 'hydrating' -%} - {{ log("Cluster " ~ cluster ~ " is hydrating: " ~ hydrated_count ~ "/" ~ total_count ~ " objects hydrated", info=True) }} - {%- elif status == 'lagging' -%} - {{ log("Cluster " ~ cluster ~ " is lagging: " ~ max_lag_secs ~ "s (threshold: " ~ threshold_secs ~ "s)", info=True) }} - {%- elif status == 'failing' -%} - {{ log("Cluster " ~ cluster ~ " is failing: " ~ failure_reason ~ " (" ~ problematic_replicas ~ "/" ~ total_replicas ~ " problematic replicas)", info=True) }} - {%- endif -%} - - {{ return({ - 'ready': status == 'ready', - 'status': status, - 'failure_reason': failure_reason, - 'hydrated_count': hydrated_count, - 'total_count': total_count, - 'max_lag_secs': max_lag_secs, - 'total_replicas': total_replicas, - 'problematic_replicas': problematic_replicas - }) }} - {%- else -%} - {{ log("Cluster " ~ cluster ~ " does not exist", info=True) }} - {{ return({ - 'ready': false, - 'status': 'failing', - 'failure_reason': 'cluster_not_found', - 'hydrated_count': 0, - 'total_count': 0, - 'max_lag_secs': 0, - 'total_replicas': 0, - 'problematic_replicas': 0 - }) }} - {%- endif -%} -{%- endif -%} +{% if execute %} + {#- are_clusters_ready reports clusters it did not find as failing, so there + is always an entry for the requested cluster. -#} + {{ return(statuses[cluster]) }} +{% endif %} {% endmacro %}