Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
27 changes: 1 addition & 26 deletions misc/dbt-materialize/dbt/adapters/materialize/impl.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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', []) %}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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) }}

Expand Down
Original file line number Diff line number Diff line change
@@ -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 %}
Original file line number Diff line number Diff line change
Expand Up @@ -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) }}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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) }}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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) %}

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

praise — I diffed the deleted SQL against are_clusters_ready and the only difference was c.name = 'x' versus c.name IN ('x'), so this really is behaviour-preserving. Worth noting for other reviewers that the cluster_not_found case is preserved too, not by this macro but by the fill-in loop at the bottom of are_clusters_ready that inserts a failing/cluster_not_found entry for every requested cluster missing from the result set. Keeping the {% if execute %} guard is also the right call, since it preserves the parse-time no-return that await_cluster_ready relies on.

written by claude on behalf of @jubrad

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for double checking the diff! The cluster_not_found fill-in loop was the part I wanted a second pair of eyes on!


{%- 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 %}
Loading