Skip to content
6 changes: 6 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -1,3 +1,9 @@
### Release [1.10.4], 2026-XX-XX

#### Bugs
* Apply column-level `codec` and `ttl` to the shard-local table of distributed models on create (with or without an enforced contract) and on add/modify column ([#741](https://github.com/ClickHouse/dbt-clickhouse/pull/741)).


### Release [1.10.3], 2026-09-15

#### Improvements
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -58,11 +58,11 @@
{% elif existing_relation.can_exchange %}
-- We can do an atomic exchange, so no need for an intermediate
{% call statement('main') -%}
{{ create_empty_table_from_relation(backup_relation, view_relation) }}
{{ create_empty_table_from_relation(backup_relation, view_relation, none, has_contract) }}
{%- endcall %}
{% do exchange_tables_atomic(backup_relation, existing_relation_local) %}
{% else %}
{% do run_query(create_empty_table_from_relation(intermediate_relation, view_relation)) or '' %}
{% do run_query(create_empty_table_from_relation(intermediate_relation, view_relation, none, has_contract)) or '' %}
{{ adapter.rename_relation(existing_relation_local, backup_relation) }}
{{ adapter.rename_relation(intermediate_relation, target_relation_local) }}
{% endif %}
Expand Down Expand Up @@ -102,31 +102,36 @@
)
{% endmacro %}

{% macro create_empty_table_from_relation(relation, source_relation, sql=none) -%}
{% macro create_empty_table_from_relation(relation, source_relation, sql=none, has_contract=false) -%}
Comment thread
cursor[bot] marked this conversation as resolved.
{%- set sql_header = config.get('sql_header', none) -%}
{%- if sql -%}
{%- set columns = adapter.get_column_schema_from_query(sql, query_settings=config.get('query_settings', {})) | list -%}
{%- else -%}
{%- set columns = adapter.get_columns_in_relation(source_relation) | list -%}
{%- endif -%}
{%- set col_list = [] -%}
{% for col in columns %}
{{col_list.append(col.name + ' ' + col.data_type) or '' }}
{% endfor %}
{{ sql_header if sql_header is not none }}

{%- if has_contract %}
{% if sql is not none %}{{ get_assert_columns_equivalent(sql) }}{% endif %}
{%- set column_defs = adapter.render_raw_columns_constraints(raw_columns=model['columns']) + adapter.render_raw_model_constraints(raw_constraints=model['constraints']) -%}
{%- else %}
{%- if sql -%}
{%- set columns = adapter.get_column_schema_from_query(sql, query_settings=config.get('query_settings', {})) | list -%}
{%- else -%}
{%- set columns = adapter.get_columns_in_relation(source_relation) | list -%}
{%- endif -%}
{%- set column_defs = [] -%}
{% for col in columns %}
{{ column_defs.append(col.name + ' ' + col.data_type + ' ' + column_codec_clause(col.name) + ' ' + column_ttl_clause(col.name)) or '' }}
Comment thread
cursor[bot] marked this conversation as resolved.
{% endfor %}
{%- endif %}

create table {{ relation.include(database=False) }}
{{ on_cluster_clause(relation) }} (
{{col_list | join(', ')}}
{{ column_defs | join(', ') }}

{% if config.get('projections') %}
{% set projections = config.get('projections') %}
{% for projection in projections %}
{% for projection in config.get('projections') %}
, {{ clickhouse_projection_ddl(projection) }}
{% endfor %}
{% endif %}
{% endif %}
)

{{ engine_clause() }}
{{ order_cols(label="order by") }}
{{ primary_key_clause(label="primary key") }}
Expand All @@ -139,7 +144,7 @@
{{ drop_relation_if_exists(shard_relation) }}
{{ drop_relation_if_exists(distributed_relation) }}
{{ create_schema(shard_relation) }}
{% do run_query(create_empty_table_from_relation(shard_relation, structure_relation, sql_query)) or '' %}
{% do run_query(create_empty_table_from_relation(shard_relation, structure_relation, sql_query, has_contract)) or '' %}
{% do run_query(create_distributed_table(distributed_relation, shard_relation)) or '' %}
{% if sql_query is not none %}
{% do run_query(clickhouse__insert_into(distributed_relation, sql_query, has_contract)) or '' %}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,13 +18,28 @@

{% endmacro %}

{% macro column_codec_clause(column_name) -%}
{{ codec_clause(model['columns'].get(column_name, {}).get('codec')) }}
{%- endmacro %}

{% macro column_ttl_clause(column_name) -%}
{{ ttl_clause(model['columns'].get(column_name, {}).get('ttl')) }}
{%- endmacro %}

{% macro exec_alter_table(relation, action, on_cluster='') %}
{% call statement('alter_table') %}
alter table {{ relation }} {{ on_cluster }} {{ action }}
{% endcall %}

{% endmacro %}

{% macro clickhouse__add_columns(columns, existing_relation, existing_local=none, is_distributed=False) %}
{% set command = 'add column if not exists' %}
{% for column in columns %}
{% set codec = model['columns'].get(column.name, {}).get('codec') %}
{% set alter_action -%}
add column if not exists `{{ column.name }}` {{ column.data_type }} {{ codec_clause(codec) }}
{%- endset %}
{% do clickhouse__run_alter_table_command(alter_action, existing_relation, existing_local, is_distributed) %}
{% set decl = '`' ~ column.name ~ '` ' ~ column.data_type %}
{% set local_action = command ~ ' ' ~ decl ~ ' ' ~ column_codec_clause(column.name) ~ ' ' ~ column_ttl_clause(column.name) %}
{% set distributed_action = command ~ ' ' ~ decl ~ ' ' ~ column_codec_clause(column.name) %}
{% do clickhouse__run_alter_table_command(local_action, existing_relation, existing_local, is_distributed, distributed_action) %}
{% endfor %}

{% endmacro %}
Expand All @@ -40,28 +55,22 @@
{% endmacro %}

{% macro clickhouse__modify_columns(columns, existing_relation, existing_local=none, is_distributed=False) %}
{% set command = 'modify column if exists' %}
{% for column in columns %}
{% set alter_action -%}
modify column if exists `{{ column.name }}` {{ column.data_type }}
{%- endset %}
{% do clickhouse__run_alter_table_command(alter_action, existing_relation, existing_local, is_distributed) %}
{% set decl = '`' ~ column.name ~ '` ' ~ column.data_type %}
{% set local_action = command ~ ' ' ~ decl ~ ' ' ~ column_codec_clause(column.name) ~ ' ' ~ column_ttl_clause(column.name) %}
{% set distributed_action = command ~ ' ' ~ decl ~ ' ' ~ column_codec_clause(column.name) %}
{% do clickhouse__run_alter_table_command(local_action, existing_relation, existing_local, is_distributed, distributed_action) %}
{% endfor %}

{% endmacro %}

{% macro clickhouse__run_alter_table_command(alter_action, existing_relation, existing_local=none, is_distributed=False) %}
{% macro clickhouse__run_alter_table_command(local_action, existing_relation, existing_local=none, is_distributed=False, distributed_action=none) %}
{% if is_distributed %}
{% call statement('alter_table') %}
alter table {{ existing_local }} {{ on_cluster_clause(existing_relation) }} {{ alter_action }}
{% endcall %}
{% call statement('alter_table') %}
alter table {{ existing_relation }} {{ on_cluster_clause(existing_relation) }} {{ alter_action }}
{% endcall %}

{% do exec_alter_table(existing_local, local_action, on_cluster_clause(existing_relation)) %}
{% do exec_alter_table(existing_relation, distributed_action or local_action, on_cluster_clause(existing_relation)) %}
{% else %}
{% call statement('alter_table') %}
alter table {{ existing_relation }} {{ alter_action }}
{% endcall %}
{% do exec_alter_table(existing_relation, local_action) %}
{% endif %}

{% endmacro %}
6 changes: 6 additions & 0 deletions dbt/include/clickhouse/macros/materializations/table.sql
Original file line number Diff line number Diff line change
Expand Up @@ -441,3 +441,9 @@
CODEC({{ codec_name }})
{%- endif %}
{% endmacro %}

{% macro ttl_clause(ttl_expr) %}
{%- if ttl_expr %}
TTL {{ ttl_expr }}
{%- endif %}
{% endmacro %}
96 changes: 80 additions & 16 deletions tests/integration/adapter/incremental/test_schema_change_codec.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,15 @@
import pytest
from dbt.tests.util import run_dbt


def assert_column_codec(project, model):
is_distributed = "distributed" in model
relation = f"{model}_local" if is_distributed else model
ddl = project.run_sql(f"SHOW CREATE TABLE {project.test_schema}.{relation}", fetch="one")[0]
assert "CODEC" in ddl
assert ("LZ4" if is_distributed else "ZSTD") in ddl


schema_change_with_codec_sql = """
{{
config(
Expand Down Expand Up @@ -84,14 +93,7 @@ def test_append_with_codec(self, project, model):
assert result[0][2] == 0
assert result[3][2] == 5

table_name = f"{project.test_schema}.{model}"
create_table_sql = project.run_sql(f"SHOW CREATE TABLE {table_name}", fetch="one")[0]

assert "CODEC" in create_table_sql
if "distributed" in model:
assert "LZ4" in create_table_sql
else:
assert "ZSTD" in create_table_sql
assert_column_codec(project, model)


sync_all_columns_with_codec_sql = """
Expand Down Expand Up @@ -166,16 +168,78 @@ def test_sync_all_columns_with_codec(self, project, model):
assert result[0][1] == 0
assert result[3][1] == 5

table_name = f"{project.test_schema}.{model}"
create_table_sql = project.run_sql(f"SHOW CREATE TABLE {table_name}", fetch="one")[0]

assert "CODEC" in create_table_sql
if "distributed" in model:
assert "LZ4" in create_table_sql
else:
assert "ZSTD" in create_table_sql
assert_column_codec(project, model)

result_types = project.run_sql(
f"select toColumnTypeName(col_1) from {model} limit 1", fetch="one"
)
assert "Float32" in result_types[0]


distributed_table_codec_sql = """
{{
config(materialized='distributed_table')
}}
select
number as col_1,
number + 1 as col_2
from numbers(3)
"""

distributed_table_codec_yml = """
version: 2
models:
- name: dist_table_rebuild_codec
config:
contract:
enforced: true
columns:
- name: col_1
data_type: UInt64
- name: col_2
data_type: UInt64
codec: LZ4
- name: dist_table_codec_no_contract
columns:
- name: col_1
data_type: UInt64
- name: col_2
data_type: UInt64
codec: LZ4
"""


class TestDistributedTableRebuildWithCodec:
@pytest.fixture(scope="class")
def models(self):
return {
"dist_table_rebuild_codec.sql": distributed_table_codec_sql,
"dist_table_codec_no_contract.sql": distributed_table_codec_sql,
"schema.yml": distributed_table_codec_yml,
}

def test_codec_survives_rebuild(self, project):
if os.environ.get('DBT_CH_TEST_CLUSTER', '').strip() == '':
pytest.skip("Not on a cluster")

run_dbt(["run", "--select", "dist_table_rebuild_codec"])
run_dbt(["run", "--select", "dist_table_rebuild_codec"])

ddl = project.run_sql(
f"SHOW CREATE TABLE {project.test_schema}.dist_table_rebuild_codec_local", fetch="one"
)[0]
assert "CODEC" in ddl
assert "LZ4" in ddl

def test_codec_applied_without_contract(self, project):
if os.environ.get('DBT_CH_TEST_CLUSTER', '').strip() == '':
pytest.skip("Not on a cluster")

run_dbt(["run", "--select", "dist_table_codec_no_contract"])

ddl = project.run_sql(
f"SHOW CREATE TABLE {project.test_schema}.dist_table_codec_no_contract_local",
fetch="one",
)[0]
assert "CODEC" in ddl
assert "LZ4" in ddl
Loading