diff --git a/CHANGELOG.md b/CHANGELOG.md index 543c0e36..6c07dfe9 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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 diff --git a/dbt/include/clickhouse/macros/materializations/distributed_table.sql b/dbt/include/clickhouse/macros/materializations/distributed_table.sql index 0e95152d..5998c282 100644 --- a/dbt/include/clickhouse/macros/materializations/distributed_table.sql +++ b/dbt/include/clickhouse/macros/materializations/distributed_table.sql @@ -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 %} @@ -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) -%} {%- 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 '' }} + {% 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") }} @@ -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 '' %} diff --git a/dbt/include/clickhouse/macros/materializations/incremental/schema_changes.sql b/dbt/include/clickhouse/macros/materializations/incremental/schema_changes.sql index f2d9080b..da9152e8 100644 --- a/dbt/include/clickhouse/macros/materializations/incremental/schema_changes.sql +++ b/dbt/include/clickhouse/macros/materializations/incremental/schema_changes.sql @@ -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 %} @@ -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 %} diff --git a/dbt/include/clickhouse/macros/materializations/table.sql b/dbt/include/clickhouse/macros/materializations/table.sql index dbff7767..b3ddf5d1 100644 --- a/dbt/include/clickhouse/macros/materializations/table.sql +++ b/dbt/include/clickhouse/macros/materializations/table.sql @@ -441,3 +441,9 @@ CODEC({{ codec_name }}) {%- endif %} {% endmacro %} + +{% macro ttl_clause(ttl_expr) %} + {%- if ttl_expr %} + TTL {{ ttl_expr }} + {%- endif %} +{% endmacro %} diff --git a/tests/integration/adapter/incremental/test_schema_change_codec.py b/tests/integration/adapter/incremental/test_schema_change_codec.py index 8c5f0a24..7f44efde 100644 --- a/tests/integration/adapter/incremental/test_schema_change_codec.py +++ b/tests/integration/adapter/incremental/test_schema_change_codec.py @@ -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( @@ -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 = """ @@ -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 diff --git a/tests/integration/adapter/incremental/test_schema_change_ttl.py b/tests/integration/adapter/incremental/test_schema_change_ttl.py new file mode 100644 index 00000000..0d700565 --- /dev/null +++ b/tests/integration/adapter/incremental/test_schema_change_ttl.py @@ -0,0 +1,245 @@ +import os + +import pytest +from dbt.tests.util import run_dbt + + +def assert_column_ttl(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 "TTL" in ddl + assert ("toIntervalDay(60)" if is_distributed else "toIntervalDay(30)") in ddl + + +schema_change_with_ttl_sql = """ +{{ + config( + materialized='%s', + unique_key='col_1', + on_schema_change='%s' + ) +}} + +{%% if not is_incremental() %%} +select + number as col_1, + number + 1 as col_2, + toDate('2020-01-01') as event_date +from numbers(3) +{%% else %%} +select + number as col_1, + number + 1 as col_2, + number + 2 as col_3, + toDate('2020-01-01') as event_date +from numbers(2, 3) +{%% endif %%} +""" + + +schema_change_with_ttl_yml = """ +version: 2 +models: + - name: schema_change_ttl_append + columns: + - name: col_1 + data_type: UInt64 + - name: col_2 + data_type: UInt64 + - name: event_date + data_type: Date + - name: col_3 + data_type: UInt64 + ttl: event_date + toIntervalDay(30) + - name: schema_change_ttl_distributed_append + columns: + - name: col_1 + data_type: UInt64 + - name: col_2 + data_type: UInt64 + - name: event_date + data_type: Date + - name: col_3 + data_type: UInt64 + ttl: event_date + toIntervalDay(60) +""" + + +class TestSchemaChangeWithTTL: + @pytest.fixture(scope="class") + def models(self): + return { + "schema_change_ttl_append.sql": schema_change_with_ttl_sql + % ("incremental", "append_new_columns"), + "schema_change_ttl_distributed_append.sql": schema_change_with_ttl_sql + % ("distributed_incremental", "append_new_columns"), + "schema.yml": schema_change_with_ttl_yml, + } + + @pytest.mark.parametrize( + "model", ("schema_change_ttl_append", "schema_change_ttl_distributed_append") + ) + def test_append_with_ttl(self, project, model): + is_distributed = "distributed" in model + if is_distributed and os.environ.get('DBT_CH_TEST_CLUSTER', '').strip() == '': + pytest.skip("Not on a cluster") + + run_dbt(["run", "--select", model]) + result = project.run_sql(f"select * from {model} order by col_1", fetch="all") + assert len(result) == 3 + + run_dbt(["--debug", "run", "--select", model]) + result = project.run_sql(f"select * from {model} order by col_1", fetch="all") + assert all(len(row) == 4 for row in result) + + assert_column_ttl(project, model) + + +sync_all_columns_with_ttl_sql = """ +{{ + config( + materialized='%s', + unique_key='col_1', + on_schema_change='sync_all_columns' + ) +}} + +{%% if not is_incremental() %%} +select + toUInt8(number) as col_1, + number + 1 as col_2, + toDate('2020-01-01') as event_date +from numbers(3) +{%% else %%} +select + toFloat32(number) as col_1, + number + 2 as col_3, + toDate('2020-01-01') as event_date +from numbers(2, 3) +{%% endif %%} +""" + +sync_all_columns_with_ttl_yml = """ +version: 2 +models: + - name: sync_ttl_test + columns: + - name: col_1 + data_type: Float32 + - name: event_date + data_type: Date + - name: col_3 + data_type: UInt64 + ttl: event_date + toIntervalDay(30) + - name: sync_ttl_distributed_test + columns: + - name: col_1 + data_type: Float32 + - name: event_date + data_type: Date + - name: col_3 + data_type: UInt64 + ttl: event_date + toIntervalDay(60) +""" + + +class TestSyncAllColumnsWithTTL: + @pytest.fixture(scope="class") + def models(self): + return { + "sync_ttl_test.sql": sync_all_columns_with_ttl_sql % "incremental", + "sync_ttl_distributed_test.sql": sync_all_columns_with_ttl_sql + % "distributed_incremental", + "schema.yml": sync_all_columns_with_ttl_yml, + } + + @pytest.mark.parametrize("model", ("sync_ttl_test", "sync_ttl_distributed_test")) + def test_sync_all_columns_with_ttl(self, project, model): + is_distributed = "distributed" in model + if is_distributed and os.environ.get('DBT_CH_TEST_CLUSTER', '').strip() == '': + pytest.skip("Not on a cluster") + + run_dbt(["run", "--select", model]) + result = project.run_sql(f"select * from {model} order by col_1", fetch="all") + assert len(result) == 3 + + run_dbt(["run", "--select", model]) + result = project.run_sql(f"select * from {model} order by col_1", fetch="all") + assert all(len(row) == 3 for row in result) + + assert_column_ttl(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_ttl_sql = """ +{{ + config(materialized='distributed_table') +}} +select + number as col_1, + toDate('2020-01-01') as event_date +from numbers(3) +""" + +distributed_table_ttl_yml = """ +version: 2 +models: + - name: dist_table_rebuild_ttl + config: + contract: + enforced: true + columns: + - name: col_1 + data_type: UInt64 + ttl: event_date + toIntervalDay(60) + - name: event_date + data_type: Date + - name: dist_table_ttl_no_contract + columns: + - name: col_1 + data_type: UInt64 + ttl: event_date + toIntervalDay(60) + - name: event_date + data_type: Date +""" + + +class TestDistributedTableRebuildWithTTL: + @pytest.fixture(scope="class") + def models(self): + return { + "dist_table_rebuild_ttl.sql": distributed_table_ttl_sql, + "dist_table_ttl_no_contract.sql": distributed_table_ttl_sql, + "schema.yml": distributed_table_ttl_yml, + } + + def test_ttl_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_ttl"]) + run_dbt(["run", "--select", "dist_table_rebuild_ttl"]) + + ddl = project.run_sql( + f"SHOW CREATE TABLE {project.test_schema}.dist_table_rebuild_ttl_local", fetch="one" + )[0] + assert "TTL" in ddl + assert "toIntervalDay(60)" in ddl + + def test_ttl_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_ttl_no_contract"]) + + ddl = project.run_sql( + f"SHOW CREATE TABLE {project.test_schema}.dist_table_ttl_no_contract_local", + fetch="one", + )[0] + assert "TTL" in ddl + assert "toIntervalDay(60)" in ddl