|
6 | 6 |
|
7 | 7 | from sqlglot import exp |
8 | 8 |
|
9 | | -from sqlmesh.core.dialect import to_schema |
| 9 | +from sqlmesh.core.dialect import to_schema, add_table |
10 | 10 | from sqlmesh.core.engine_adapter.base import ( |
11 | 11 | EngineAdapterWithIndexSupport, |
12 | 12 | EngineAdapter, |
13 | 13 | InsertOverwriteStrategy, |
| 14 | + MERGE_SOURCE_ALIAS, |
| 15 | + MERGE_TARGET_ALIAS, |
14 | 16 | ) |
15 | 17 | from sqlmesh.core.engine_adapter.mixins import ( |
16 | 18 | GetCurrentCatalogFromFunctionMixin, |
|
32 | 34 |
|
33 | 35 | if t.TYPE_CHECKING: |
34 | 36 | from sqlmesh.core._typing import SchemaName, TableName |
35 | | - from sqlmesh.core.engine_adapter._typing import DF, Query |
| 37 | + from sqlmesh.core.engine_adapter._typing import DF, Query, QueryOrDF |
36 | 38 |
|
37 | 39 |
|
38 | 40 | @set_catalog() |
@@ -188,6 +190,87 @@ def drop_schema( |
188 | 190 | ) |
189 | 191 | super().drop_schema(schema_name, ignore_if_not_exists=ignore_if_not_exists, cascade=False) |
190 | 192 |
|
| 193 | + def merge( |
| 194 | + self, |
| 195 | + target_table: TableName, |
| 196 | + source_table: QueryOrDF, |
| 197 | + columns_to_types: t.Optional[t.Dict[str, exp.DataType]], |
| 198 | + unique_key: t.Sequence[exp.Expression], |
| 199 | + when_matched: t.Optional[exp.Whens] = None, |
| 200 | + merge_filter: t.Optional[exp.Expression] = None, |
| 201 | + ) -> None: |
| 202 | + source_queries, columns_to_types = self._get_source_queries_and_columns_to_types( |
| 203 | + source_table, columns_to_types, target_table=target_table |
| 204 | + ) |
| 205 | + columns_to_types = columns_to_types or self.columns(target_table) |
| 206 | + on = exp.and_( |
| 207 | + *( |
| 208 | + add_table(part, MERGE_TARGET_ALIAS).eq(add_table(part, MERGE_SOURCE_ALIAS)) |
| 209 | + for part in unique_key |
| 210 | + ) |
| 211 | + ) |
| 212 | + if merge_filter: |
| 213 | + on = exp.and_(merge_filter, on) |
| 214 | + |
| 215 | + if not when_matched: |
| 216 | + match_condition = None |
| 217 | + unique_key_names = [y.name for y in unique_key] |
| 218 | + columns_to_types_no_keys = [c for c in columns_to_types if c not in unique_key_names] |
| 219 | + |
| 220 | + target_columns_no_keys = [ |
| 221 | + exp.column(c, MERGE_TARGET_ALIAS) for c in columns_to_types_no_keys |
| 222 | + ] |
| 223 | + source_columns_no_keys = [ |
| 224 | + exp.column(c, MERGE_SOURCE_ALIAS) for c in columns_to_types_no_keys |
| 225 | + ] |
| 226 | + |
| 227 | + match_condition = exp.Exists( |
| 228 | + this=exp.select(*target_columns_no_keys).except_( |
| 229 | + exp.select(*source_columns_no_keys) |
| 230 | + ) |
| 231 | + ) |
| 232 | + |
| 233 | + match_expressions = [ |
| 234 | + exp.When( |
| 235 | + matched=True, |
| 236 | + source=False, |
| 237 | + condition=match_condition, |
| 238 | + then=exp.Update( |
| 239 | + expressions=[ |
| 240 | + exp.column(col, MERGE_TARGET_ALIAS).eq( |
| 241 | + exp.column(col, MERGE_SOURCE_ALIAS) |
| 242 | + ) |
| 243 | + for col in columns_to_types_no_keys |
| 244 | + ], |
| 245 | + ), |
| 246 | + ) |
| 247 | + ] |
| 248 | + else: |
| 249 | + match_expressions = when_matched.copy().expressions |
| 250 | + |
| 251 | + match_expressions.append( |
| 252 | + exp.When( |
| 253 | + matched=False, |
| 254 | + source=False, |
| 255 | + then=exp.Insert( |
| 256 | + this=exp.Tuple(expressions=[exp.column(col) for col in columns_to_types]), |
| 257 | + expression=exp.Tuple( |
| 258 | + expressions=[ |
| 259 | + exp.column(col, MERGE_SOURCE_ALIAS) for col in columns_to_types |
| 260 | + ] |
| 261 | + ), |
| 262 | + ), |
| 263 | + ) |
| 264 | + ) |
| 265 | + for source_query in source_queries: |
| 266 | + with source_query as query: |
| 267 | + self._merge( |
| 268 | + target_table=target_table, |
| 269 | + query=query, |
| 270 | + on=on, |
| 271 | + whens=exp.Whens(expressions=match_expressions), |
| 272 | + ) |
| 273 | + |
191 | 274 | def _convert_df_datetime(self, df: DF, columns_to_types: t.Dict[str, exp.DataType]) -> None: |
192 | 275 | import pandas as pd |
193 | 276 | from pandas.api.types import is_datetime64_any_dtype # type: ignore |
|
0 commit comments