Skip to content

Commit b16841e

Browse files
authored
Add generalized (relations …) CSV loading construct (#259)
1 parent 766469e commit b16841e

31 files changed

Lines changed: 16874 additions & 14156 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
@@ -1107,13 +1115,58 @@ csv_asof
11071115
: "(" "asof" STRING ")"
11081116

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

11181171
csv_locator_paths
11191172
: "(" "paths" STRING* ")"
@@ -1473,6 +1526,62 @@ def _try_extract_value_string_list(value: Optional[logic.Value]) -> Optional[Seq
14731526
return None
14741527

14751528

1529+
def construct_non_cdc_relations(
1530+
targets: Sequence[logic.TargetRelation],
1531+
) -> logic.TargetRelations:
1532+
return logic.TargetRelations(
1533+
keys=list[logic.NamedColumn](),
1534+
plain=logic.PlainTargets(targets=targets),
1535+
)
1536+
1537+
1538+
def construct_cdc_relations(
1539+
inserts: Sequence[logic.TargetRelation],
1540+
deletes: Sequence[logic.TargetRelation],
1541+
) -> logic.TargetRelations:
1542+
return logic.TargetRelations(
1543+
keys=list[logic.NamedColumn](),
1544+
cdc=logic.CDCTargets(inserts=inserts, deletes=deletes),
1545+
)
1546+
1547+
1548+
def construct_relations(
1549+
keys: Sequence[logic.NamedColumn],
1550+
body: logic.TargetRelations,
1551+
) -> logic.TargetRelations:
1552+
if builtin.has_proto_field(body, "plain"):
1553+
return logic.TargetRelations(keys=keys, plain=body.plain)
1554+
return logic.TargetRelations(keys=keys, cdc=body.cdc)
1555+
1556+
1557+
def construct_csv_data(
1558+
locator: logic.CSVLocator,
1559+
config: logic.CSVConfig,
1560+
columns_opt: Optional[Sequence[logic.GNFColumn]],
1561+
relations_opt: Optional[logic.TargetRelations],
1562+
asof: String,
1563+
) -> logic.CSVData:
1564+
return logic.CSVData(
1565+
locator=locator,
1566+
config=config,
1567+
columns=builtin.unwrap_option_or(columns_opt, list[logic.GNFColumn]()),
1568+
asof=asof,
1569+
relations=relations_opt,
1570+
)
1571+
1572+
1573+
def deconstruct_csv_data_columns_optional(msg: logic.CSVData) -> Optional[Sequence[logic.GNFColumn]]:
1574+
if builtin.has_proto_field(msg, "relations"):
1575+
return builtin.none()
1576+
return builtin.some(msg.columns)
1577+
1578+
1579+
def deconstruct_csv_data_relations_optional(msg: logic.CSVData) -> Optional[logic.TargetRelations]:
1580+
if builtin.has_proto_field(msg, "relations"):
1581+
return builtin.some(builtin.unwrap_option(msg.relations))
1582+
return builtin.none()
1583+
1584+
14761585
def construct_csv_config(
14771586
config_dict: Sequence[Tuple[String, logic.Value]],
14781587
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)