Skip to content

Commit 912ec50

Browse files
committed
Merge remote-tracking branch 'origin/main' into dz-export-output-2
2 parents e9c69c9 + 2a0419a commit 912ec50

33 files changed

Lines changed: 25877 additions & 7409 deletions

meta/src/meta/grammar.y

Lines changed: 113 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -87,6 +87,14 @@
8787
%nonterm gnf_column logic.GNFColumn
8888
%nonterm gnf_column_path Sequence[String]
8989
%nonterm gnf_columns Sequence[logic.GNFColumn]
90+
%nonterm named_column logic.NamedColumn
91+
%nonterm relation_keys Sequence[logic.NamedColumn]
92+
%nonterm target_relation logic.TargetRelation
93+
%nonterm non_cdc_relations Sequence[logic.TargetRelation]
94+
%nonterm cdc_inserts Sequence[logic.TargetRelation]
95+
%nonterm cdc_deletes Sequence[logic.TargetRelation]
96+
%nonterm relation_body logic.TargetRelations
97+
%nonterm target_relations logic.TargetRelations
9098
%nonterm csv_config logic.CSVConfig
9199
%nonterm csv_data logic.CSVData
92100
%nonterm csv_locator_inline_data String
@@ -1108,13 +1116,58 @@ csv_asof
11081116
: "(" "asof" STRING ")"
11091117

11101118
csv_data
1111-
: "(" "csv_data" csvlocator csv_config gnf_columns csv_asof ")"
1112-
construct: $$ = logic.CSVData(locator=$3, config=$4, columns=$5, asof=$6)
1119+
: "(" "csv_data" csvlocator csv_config gnf_columns? target_relations? csv_asof ")"
1120+
construct: $$ = construct_csv_data($3, $4, $5, $6, $7)
11131121
deconstruct:
11141122
$3: logic.CSVLocator = $$.locator
11151123
$4: logic.CSVConfig = $$.config
1116-
$5: Sequence[logic.GNFColumn] = $$.columns
1117-
$6: String = $$.asof
1124+
$5: Optional[Sequence[logic.GNFColumn]] = deconstruct_csv_data_columns_optional($$)
1125+
$6: Optional[logic.TargetRelations] = deconstruct_csv_data_relations_optional($$)
1126+
$7: String = $$.asof
1127+
1128+
named_column
1129+
: "(" "column" STRING type ")"
1130+
construct: $$ = logic.NamedColumn(name=$3, type=$4)
1131+
deconstruct:
1132+
$3: String = $$.name
1133+
$4: logic.Type = $$.type
1134+
1135+
relation_keys
1136+
: "(" "keys" named_column* ")"
1137+
1138+
target_relation
1139+
: "(" "relation" relation_id named_column* ")"
1140+
construct: $$ = logic.TargetRelation(target_id=$3, values=$4)
1141+
deconstruct:
1142+
$3: logic.RelationId = $$.target_id
1143+
$4: Sequence[logic.NamedColumn] = $$.values
1144+
1145+
non_cdc_relations
1146+
: target_relation*
1147+
1148+
cdc_inserts
1149+
: "(" "inserts" target_relation* ")"
1150+
1151+
cdc_deletes
1152+
: "(" "deletes" target_relation* ")"
1153+
1154+
relation_body
1155+
: non_cdc_relations
1156+
construct: $$ = construct_non_cdc_relations($1)
1157+
deconstruct if builtin.has_proto_field($$, 'plain'):
1158+
$1: Sequence[logic.TargetRelation] = $$.plain.targets
1159+
| cdc_inserts cdc_deletes
1160+
construct: $$ = construct_cdc_relations($1, $2)
1161+
deconstruct if builtin.has_proto_field($$, 'cdc'):
1162+
$1: Sequence[logic.TargetRelation] = $$.cdc.inserts
1163+
$2: Sequence[logic.TargetRelation] = $$.cdc.deletes
1164+
1165+
target_relations
1166+
: "(" "relations" relation_keys relation_body ")"
1167+
construct: $$ = construct_relations($3, $4)
1168+
deconstruct:
1169+
$3: Sequence[logic.NamedColumn] = $$.keys
1170+
$4: logic.TargetRelations = $$
11181171

11191172
csv_locator_paths
11201173
: "(" "paths" STRING* ")"
@@ -1484,6 +1537,62 @@ def _try_extract_value_string_list(value: Optional[logic.Value]) -> Optional[Seq
14841537
return None
14851538

14861539

1540+
def construct_non_cdc_relations(
1541+
targets: Sequence[logic.TargetRelation],
1542+
) -> logic.TargetRelations:
1543+
return logic.TargetRelations(
1544+
keys=list[logic.NamedColumn](),
1545+
plain=logic.PlainTargets(targets=targets),
1546+
)
1547+
1548+
1549+
def construct_cdc_relations(
1550+
inserts: Sequence[logic.TargetRelation],
1551+
deletes: Sequence[logic.TargetRelation],
1552+
) -> logic.TargetRelations:
1553+
return logic.TargetRelations(
1554+
keys=list[logic.NamedColumn](),
1555+
cdc=logic.CDCTargets(inserts=inserts, deletes=deletes),
1556+
)
1557+
1558+
1559+
def construct_relations(
1560+
keys: Sequence[logic.NamedColumn],
1561+
body: logic.TargetRelations,
1562+
) -> logic.TargetRelations:
1563+
if builtin.has_proto_field(body, "plain"):
1564+
return logic.TargetRelations(keys=keys, plain=body.plain)
1565+
return logic.TargetRelations(keys=keys, cdc=body.cdc)
1566+
1567+
1568+
def construct_csv_data(
1569+
locator: logic.CSVLocator,
1570+
config: logic.CSVConfig,
1571+
columns_opt: Optional[Sequence[logic.GNFColumn]],
1572+
relations_opt: Optional[logic.TargetRelations],
1573+
asof: String,
1574+
) -> logic.CSVData:
1575+
return logic.CSVData(
1576+
locator=locator,
1577+
config=config,
1578+
columns=builtin.unwrap_option_or(columns_opt, list[logic.GNFColumn]()),
1579+
asof=asof,
1580+
relations=relations_opt,
1581+
)
1582+
1583+
1584+
def deconstruct_csv_data_columns_optional(msg: logic.CSVData) -> Optional[Sequence[logic.GNFColumn]]:
1585+
if builtin.has_proto_field(msg, "relations"):
1586+
return builtin.none()
1587+
return builtin.some(msg.columns)
1588+
1589+
1590+
def deconstruct_csv_data_relations_optional(msg: logic.CSVData) -> Optional[logic.TargetRelations]:
1591+
if builtin.has_proto_field(msg, "relations"):
1592+
return builtin.some(builtin.unwrap_option(msg.relations))
1593+
return builtin.none()
1594+
1595+
14871596
def construct_csv_config(
14881597
config_dict: Sequence[Tuple[String, logic.Value]],
14891598
storage_integration_opt: Optional[Sequence[Tuple[String, logic.Value]]],

proto/relationalai/lqp/v1/logic.proto

Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -288,11 +288,46 @@ message StorageIntegration {
288288
string s3_secret_access_key = 5;
289289
}
290290

291+
// A single named column with its type. Used to describe both shared key columns and
292+
// per-relation value columns in the generalized `TargetRelations` loading construct.
293+
message NamedColumn {
294+
string name = 1; // Column name (e.g. "src")
295+
Type type = 2; // Column type
296+
}
297+
298+
// One target relation: the shared keys plus this relation's own (possibly empty) value columns.
299+
message TargetRelation {
300+
RelationId target_id = 1; // Output relation path
301+
repeated NamedColumn values = 2; // Value columns for this relation (may be empty)
302+
}
303+
304+
// Plain (non-CDC) load: each input row becomes a tuple in every target relation.
305+
message PlainTargets {
306+
repeated TargetRelation targets = 1; // Target relations (each row feeds all of them)
307+
}
308+
309+
// CDC load: input rows are routed by METADATA$ACTION into insert and delete deltas.
310+
message CDCTargets {
311+
repeated TargetRelation inserts = 1; // INSERT-action rows feed these
312+
repeated TargetRelation deletes = 2; // DELETE-action rows feed these
313+
}
314+
315+
// Generalized loading: shared key columns plus the target relations, loaded either as a
316+
// plain snapshot or as CDC insert/delete deltas. The two modes are mutually exclusive.
317+
message TargetRelations {
318+
repeated NamedColumn keys = 1; // Shared key columns
319+
oneof body {
320+
PlainTargets plain = 2;
321+
CDCTargets cdc = 3;
322+
}
323+
}
324+
291325
message CSVData {
292326
CSVLocator locator = 1;
293327
CSVConfig config = 2;
294328
repeated GNFColumn columns = 3;
295329
string asof = 4; // Blob storage timestamp for freshness requirements
330+
optional TargetRelations relations = 5; // If present, generalized loading; mutually exclusive with columns
296331
}
297332

298333
message CSVLocator {

0 commit comments

Comments
 (0)