Skip to content

Commit 5ef3a92

Browse files
authored
Merge branch 'main' into fix/snowflake-iceberg-column-comments
2 parents 24a4083 + 2c30f83 commit 5ef3a92

10 files changed

Lines changed: 108 additions & 11 deletions

File tree

‎sqlmesh/core/engine_adapter/base.py‎

Lines changed: 13 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1098,6 +1098,8 @@ def clone_table(
10981098
replace: bool = False,
10991099
exists: bool = True,
11001100
clone_kwargs: t.Optional[t.Dict[str, t.Any]] = None,
1101+
table_format: t.Optional[str] = None,
1102+
table_kind: t.Optional[str] = None,
11011103
**kwargs: t.Any,
11021104
) -> None:
11031105
"""Creates a table with the target name by cloning the source table.
@@ -1107,6 +1109,10 @@ def clone_table(
11071109
source_table_name: The name of the source table that should be cloned.
11081110
replace: Whether or not to replace an existing table.
11091111
exists: Indicates whether to include the IF NOT EXISTS check.
1112+
clone_kwargs: Additional arguments for the CLONE clause.
1113+
table_format: The table format of the source table, if any. Engines that require
1114+
format-specific DDL to clone a table use it to derive `table_kind`.
1115+
table_kind: The kind of table to create. Defaults to `TABLE`.
11101116
"""
11111117
if not self.SUPPORTS_CLONING:
11121118
raise NotImplementedError(f"Engine does not support cloning: {type(self)}")
@@ -1115,7 +1121,7 @@ def clone_table(
11151121
self.execute(
11161122
exp.Create(
11171123
this=exp.to_table(target_table_name),
1118-
kind="TABLE",
1124+
kind=table_kind or "TABLE",
11191125
replace=replace,
11201126
exists=exists,
11211127
clone=exp.Clone(
@@ -1218,9 +1224,15 @@ def get_alter_operations(
12181224
def alter_table(
12191225
self,
12201226
alter_expressions: t.Union[t.List[exp.Alter], t.List[TableAlterOperation]],
1227+
table_format: t.Optional[str] = None,
12211228
) -> None:
12221229
"""
12231230
Performs the alter statements to change the current table into the structure of the target table.
1231+
1232+
Args:
1233+
alter_expressions: The alter operations to apply.
1234+
table_format: The table format of the target table, if any. Engines that require
1235+
format-specific DDL to alter a table use it to adjust the generated statements.
12241236
"""
12251237
with self.transaction():
12261238
for alter_expression in [

‎sqlmesh/core/engine_adapter/bigquery.py‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -405,6 +405,7 @@ def create_mapping_schema(
405405
def alter_table(
406406
self,
407407
alter_expressions: t.Union[t.List[exp.Alter], t.List[TableAlterOperation]],
408+
table_format: t.Optional[str] = None,
408409
) -> None:
409410
"""
410411
Performs the alter statements to change the current table into the structure of the target table,

‎sqlmesh/core/engine_adapter/clickhouse.py‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -699,6 +699,7 @@ def delete_from(self, table_name: TableName, where: t.Union[str, exp.Expr]) -> N
699699
def alter_table(
700700
self,
701701
alter_expressions: t.Union[t.List[exp.Alter], t.List[TableAlterOperation]],
702+
table_format: t.Optional[str] = None,
702703
) -> None:
703704
"""
704705
Performs the alter statements to change the current table into the structure of the target table.

‎sqlmesh/core/engine_adapter/databricks.py‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -386,6 +386,8 @@ def clone_table(
386386
replace: bool = False,
387387
exists: bool = True,
388388
clone_kwargs: t.Optional[t.Dict[str, t.Any]] = None,
389+
table_format: t.Optional[str] = None,
390+
table_kind: t.Optional[str] = None,
389391
**kwargs: t.Any,
390392
) -> None:
391393
clone_kwargs = clone_kwargs or {}
@@ -395,6 +397,8 @@ def clone_table(
395397
source_table_name,
396398
replace=replace,
397399
clone_kwargs=clone_kwargs,
400+
table_format=table_format,
401+
table_kind=table_kind,
398402
**kwargs,
399403
)
400404

‎sqlmesh/core/engine_adapter/fabric.py‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -225,7 +225,9 @@ def set_current_catalog(self, catalog_name: t.Optional[str]) -> None:
225225
self._target_catalog = target_catalog
226226

227227
def alter_table(
228-
self, alter_expressions: t.Union[t.List[exp.Alter], t.List[TableAlterOperation]]
228+
self,
229+
alter_expressions: t.Union[t.List[exp.Alter], t.List[TableAlterOperation]],
230+
table_format: t.Optional[str] = None,
229231
) -> None:
230232
"""
231233
Applies alter expressions to a table. Fabric has limited support for ALTER TABLE,

‎sqlmesh/core/engine_adapter/snowflake.py‎

Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,7 @@
2424
SourceQuery,
2525
set_catalog,
2626
)
27+
from sqlmesh.core.schema_diff import TableAlterOperation
2728
from sqlmesh.utils import optional_import, get_source_columns_to_types
2829
from sqlmesh.utils.errors import SQLMeshError
2930
from sqlmesh.utils.pandas import columns_to_types_from_dtypes
@@ -674,6 +675,8 @@ def clone_table(
674675
replace: bool = False,
675676
exists: bool = True,
676677
clone_kwargs: t.Optional[t.Dict[str, t.Any]] = None,
678+
table_format: t.Optional[str] = None,
679+
table_kind: t.Optional[str] = None,
677680
**kwargs: t.Any,
678681
) -> None:
679682
# The Snowflake adapter should use the transient property to clone transient tables
@@ -682,14 +685,43 @@ def clone_table(
682685
if isinstance(table_type, exp.TransientProperty):
683686
kwargs["properties"] = exp.Properties(expressions=[table_type])
684687

688+
# Snowflake rejects `CREATE TABLE ... CLONE` for Iceberg tables, it requires
689+
# `CREATE ICEBERG TABLE ... CLONE` instead
690+
if table_format and not table_kind:
691+
table_kind = f"{table_format.upper()} TABLE"
692+
685693
super().clone_table(
686694
target_table_name,
687695
source_table_name,
688696
replace=replace,
689697
clone_kwargs=clone_kwargs,
698+
table_kind=table_kind,
690699
**kwargs,
691700
)
692701

702+
def alter_table(
703+
self,
704+
alter_expressions: t.Union[t.List[exp.Alter], t.List[TableAlterOperation]],
705+
table_format: t.Optional[str] = None,
706+
) -> None:
707+
# Snowflake rejects `ALTER TABLE` for Iceberg tables, it requires
708+
# `ALTER ICEBERG TABLE` instead
709+
if table_format:
710+
table_kind = f"{table_format.upper()} TABLE"
711+
resolved_expressions = []
712+
for alter_expression in alter_expressions:
713+
resolved_expression = (
714+
alter_expression.expression
715+
if isinstance(alter_expression, TableAlterOperation)
716+
else alter_expression.copy()
717+
)
718+
resolved_expression.set("kind", table_kind)
719+
resolved_expressions.append(resolved_expression)
720+
721+
super().alter_table(resolved_expressions)
722+
else:
723+
super().alter_table(alter_expressions)
724+
693725
@t.overload
694726
def _columns_to_types(
695727
self,

‎sqlmesh/core/snapshot/evaluator.py‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1104,6 +1104,7 @@ def _clone_snapshot_in_dev(
11041104
target_table_name,
11051105
snapshot.table_name(),
11061106
rendered_physical_properties=rendered_physical_properties,
1107+
table_format=snapshot.model.table_format,
11071108
)
11081109
self._migrate_target_table(
11091110
target_table_name=target_table_name,
@@ -2161,7 +2162,7 @@ def migrate(
21612162
_check_additive_schema_change(
21622163
snapshot, alter_operations, kwargs["allow_additive_snapshots"]
21632164
)
2164-
self.adapter.alter_table(alter_operations)
2165+
self.adapter.alter_table(alter_operations, table_format=snapshot.model.table_format)
21652166

21662167
# Apply grants after schema migration
21672168
deployability_index = kwargs.get("deployability_index")

‎tests/core/engine_adapter/integration/docker/_common-hive.yaml‎

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -12,22 +12,22 @@ services:
1212

1313
# S3-style object storage
1414
minio:
15-
image: 'minio/minio:RELEASE.2022-05-26T05-48-41Z'
15+
image: 'cgr.dev/chainguard/minio:latest@sha256:039800e64ec7247d2fde7cff3697e964f6fe20b6d7d2c46aa7d82cc63355d512'
1616
ports:
1717
- '9000:9000'
1818
- '9001:9001'
1919
environment:
20-
MINIO_ACCESS_KEY: minio
21-
MINIO_SECRET_KEY: minio123
20+
MINIO_ROOT_USER: minio
21+
MINIO_ROOT_PASSWORD: minio123
2222
command: server /data --console-address ":9001"
2323

2424
# Set up minio with default buckets / paths
2525
mc-job:
26-
image: 'minio/mc:RELEASE.2022-05-09T04-08-26Z'
26+
image: 'cgr.dev/chainguard/minio-client:latest-dev@sha256:fb635b967f5f32424391150dac4985cc1aea78fb0a1ce031c01952cf661a13ec'
2727
entrypoint: |
2828
/bin/bash -c "
2929
sleep 5;
30-
/usr/bin/mc config --quiet host add myminio http://minio:9000 minio minio123;
30+
/usr/bin/mc alias --quiet set myminio http://minio:9000 minio minio123;
3131
/usr/bin/mc mb --quiet myminio/trino/datalake;
3232
/usr/bin/mc mb --quiet myminio/trino/datalake_iceberg;
3333
/usr/bin/mc mb --quiet myminio/trino/datalake_delta;
@@ -39,4 +39,4 @@ services:
3939
/usr/bin/mc mb --quiet myminio/nessie/warehouse;
4040
"
4141
depends_on:
42-
- minio
42+
- minio

‎tests/core/engine_adapter/test_snowflake.py‎

Lines changed: 42 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1039,6 +1039,48 @@ def test_table_format_iceberg(snowflake_mocked_engine_adapter: SnowflakeEngineAd
10391039
]
10401040

10411041

1042+
def test_clone_table_iceberg(mocker: MockerFixture, make_mocked_engine_adapter: t.Callable):
1043+
mocker.patch("sqlmesh.core.engine_adapter.snowflake.SnowflakeEngineAdapter.set_current_catalog")
1044+
adapter = make_mocked_engine_adapter(SnowflakeEngineAdapter, default_catalog="test_catalog")
1045+
1046+
# Snowflake rejects `CREATE TABLE ... CLONE` for Iceberg tables
1047+
adapter.clone_table("target_table", "source_table", table_format="iceberg")
1048+
adapter.cursor.execute.assert_called_once_with(
1049+
'CREATE ICEBERG TABLE IF NOT EXISTS "target_table" CLONE "source_table"'
1050+
)
1051+
1052+
# Engines that don't need format-specific DDL are unaffected
1053+
adapter = make_mocked_engine_adapter(EngineAdapter, default_catalog="test_catalog")
1054+
adapter.SUPPORTS_CLONING = True
1055+
adapter.clone_table("target_table", "source_table", table_format="iceberg")
1056+
adapter.cursor.execute.assert_called_once_with(
1057+
'CREATE TABLE IF NOT EXISTS "target_table" CLONE "source_table"'
1058+
)
1059+
1060+
1061+
def test_alter_table_iceberg(mocker: MockerFixture, make_mocked_engine_adapter: t.Callable):
1062+
mocker.patch("sqlmesh.core.engine_adapter.snowflake.SnowflakeEngineAdapter.set_current_catalog")
1063+
adapter = make_mocked_engine_adapter(SnowflakeEngineAdapter, default_catalog="test_catalog")
1064+
1065+
current_table = {"a": "INT"}
1066+
target_table = {"a": "INT", "b": "INT"}
1067+
adapter.columns = lambda table_name, **kwargs: {
1068+
k: exp.DataType.build(v)
1069+
for k, v in (current_table if table_name == "test_table" else target_table).items()
1070+
}
1071+
1072+
alter_operations = adapter.get_alter_operations("test_table", "target_table")
1073+
1074+
# Snowflake rejects `ALTER TABLE` for Iceberg tables
1075+
adapter.alter_table(alter_operations, table_format="iceberg")
1076+
assert to_sql_calls(adapter) == ['ALTER ICEBERG TABLE "test_table" ADD "b" INT']
1077+
1078+
# Without a table format the regular `ALTER TABLE` is used
1079+
adapter = make_mocked_engine_adapter(SnowflakeEngineAdapter, default_catalog="test_catalog")
1080+
adapter.alter_table(alter_operations)
1081+
assert to_sql_calls(adapter) == ['ALTER TABLE "test_table" ADD "b" INT']
1082+
1083+
10421084
def test_create_view_with_schema_and_grants(
10431085
snowflake_mocked_engine_adapter: SnowflakeEngineAdapter,
10441086
):

‎tests/core/test_snapshot_evaluator.py‎

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1943,6 +1943,7 @@ def test_create_clone_in_dev(mocker: MockerFixture, adapter_mock, make_snapshot)
19431943
f"sqlmesh__test_schema.test_schema__test_model__{snapshot.dev_version}__dev",
19441944
f"sqlmesh__test_schema.test_schema__test_model__{snapshot.version}",
19451945
rendered_physical_properties={},
1946+
table_format=None,
19461947
)
19471948

19481949
adapter_mock.get_alter_operations.assert_called_once_with(
@@ -1952,7 +1953,7 @@ def test_create_clone_in_dev(mocker: MockerFixture, adapter_mock, make_snapshot)
19521953
ignore_additive=False,
19531954
)
19541955

1955-
adapter_mock.alter_table.assert_called_once_with([])
1956+
adapter_mock.alter_table.assert_called_once_with([], table_format=None)
19561957

19571958
adapter_mock.drop_table.assert_called_once_with(
19581959
f"sqlmesh__test_schema.test_schema__test_model__{snapshot.version}__dev_schema_tmp"
@@ -1992,6 +1993,7 @@ def test_drop_clone_in_dev_when_migration_fails(mocker: MockerFixture, adapter_m
19921993
f"sqlmesh__test_schema.test_schema__test_model__{snapshot.version}__dev",
19931994
f"sqlmesh__test_schema.test_schema__test_model__{snapshot.version}",
19941995
rendered_physical_properties={},
1996+
table_format=None,
19951997
)
19961998

19971999
adapter_mock.get_alter_operations.assert_called_once_with(
@@ -2001,7 +2003,7 @@ def test_drop_clone_in_dev_when_migration_fails(mocker: MockerFixture, adapter_m
20012003
ignore_additive=False,
20022004
)
20032005

2004-
adapter_mock.alter_table.assert_called_once_with([])
2006+
adapter_mock.alter_table.assert_called_once_with([], table_format=None)
20052007

20062008
adapter_mock.drop_table.assert_has_calls(
20072009
[

0 commit comments

Comments
 (0)