Skip to content

Commit 3335293

Browse files
hbarthelsclaude
andcommitted
Add generalized (relations …) CSV loading construct
Adds a new `(relations …)` construct on `CSVData` alongside the legacy `(columns …)` form: a shared set of key columns (or the special `METADATA$KEY`) plus one or more output relations, each with its own (possibly empty) value columns, with optional CDC `(inserts …)`/`(deletes …)` grouping. - proto: `NamedColumn` / `OutputRelation` / `Relations` messages + an optional `Relations relations` field on `CSVData` (mutually exclusive with `columns`) - grammar: `relations` / `relation_keys` / `output_relation` / `named_column` rules; `relation_body` returns a concrete `Relations` to keep the Go parser type-stable - regenerated Python / Julia / Go parsers, pretty-printers, and protobuf bindings; `global_ids` + equality extended for the new messages (Julia SDK) - `.lqp` fixtures (binary edge, arity-4 edge, two-relation split, CDC) + regenerated bin / pretty / pretty_debug snapshots `make test` green across Python, Julia, and Go. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
1 parent 00da0f1 commit 3335293

30 files changed

Lines changed: 16277 additions & 14159 deletions

meta/src/meta/grammar.y

Lines changed: 120 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 output_relation logic.OutputRelation
93+
%nonterm non_cdc_relations Sequence[logic.OutputRelation]
94+
%nonterm cdc_inserts Sequence[logic.OutputRelation]
95+
%nonterm cdc_deletes Sequence[logic.OutputRelation]
96+
%nonterm relation_body logic.Relations
97+
%nonterm relations logic.Relations
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? 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.Relations] = 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+
output_relation
1138+
: "(" "relation" relation_id named_column* ")"
1139+
construct: $$ = logic.OutputRelation(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+
: output_relation*
1146+
1147+
cdc_inserts
1148+
: "(" "inserts" output_relation* ")"
1149+
1150+
cdc_deletes
1151+
: "(" "deletes" output_relation* ")"
1152+
1153+
relation_body
1154+
: non_cdc_relations
1155+
construct: $$ = construct_non_cdc_relations($1)
1156+
deconstruct if builtin.is_empty($$.inserts) and builtin.is_empty($$.deletes):
1157+
$1: Sequence[logic.OutputRelation] = $$.relations
1158+
| cdc_inserts cdc_deletes
1159+
construct: $$ = construct_cdc_relations($1, $2)
1160+
deconstruct if not (builtin.is_empty($$.inserts) and builtin.is_empty($$.deletes)):
1161+
$1: Sequence[logic.OutputRelation] = $$.inserts
1162+
$2: Sequence[logic.OutputRelation] = $$.deletes
1163+
1164+
relations
1165+
: "(" "relations" relation_keys relation_body ")"
1166+
construct: $$ = construct_relations($3, $4)
1167+
deconstruct:
1168+
$3: Sequence[logic.NamedColumn] = $$.keys
1169+
$4: logic.Relations = $$
11171170

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

14751528

1529+
def construct_non_cdc_relations(
1530+
relations: Sequence[logic.OutputRelation],
1531+
) -> logic.Relations:
1532+
return logic.Relations(
1533+
keys=list[logic.NamedColumn](),
1534+
relations=relations,
1535+
inserts=list[logic.OutputRelation](),
1536+
deletes=list[logic.OutputRelation](),
1537+
)
1538+
1539+
1540+
def construct_cdc_relations(
1541+
inserts: Sequence[logic.OutputRelation],
1542+
deletes: Sequence[logic.OutputRelation],
1543+
) -> logic.Relations:
1544+
return logic.Relations(
1545+
keys=list[logic.NamedColumn](),
1546+
relations=list[logic.OutputRelation](),
1547+
inserts=inserts,
1548+
deletes=deletes,
1549+
)
1550+
1551+
1552+
def construct_relations(
1553+
keys: Sequence[logic.NamedColumn],
1554+
body: logic.Relations,
1555+
) -> logic.Relations:
1556+
return logic.Relations(
1557+
keys=keys,
1558+
relations=body.relations,
1559+
inserts=body.inserts,
1560+
deletes=body.deletes,
1561+
)
1562+
1563+
1564+
def construct_csv_data(
1565+
locator: logic.CSVLocator,
1566+
config: logic.CSVConfig,
1567+
columns_opt: Optional[Sequence[logic.GNFColumn]],
1568+
relations_opt: Optional[logic.Relations],
1569+
asof: String,
1570+
) -> logic.CSVData:
1571+
return logic.CSVData(
1572+
locator=locator,
1573+
config=config,
1574+
columns=builtin.unwrap_option_or(columns_opt, list[logic.GNFColumn]()),
1575+
asof=asof,
1576+
relations=relations_opt,
1577+
)
1578+
1579+
1580+
def deconstruct_csv_data_columns_optional(msg: logic.CSVData) -> Optional[Sequence[logic.GNFColumn]]:
1581+
if builtin.has_proto_field(msg, "relations"):
1582+
return builtin.none()
1583+
return builtin.some(msg.columns)
1584+
1585+
1586+
def deconstruct_csv_data_relations_optional(msg: logic.CSVData) -> Optional[logic.Relations]:
1587+
if builtin.has_proto_field(msg, "relations"):
1588+
return builtin.some(builtin.unwrap_option(msg.relations))
1589+
return builtin.none()
1590+
1591+
14761592
def construct_csv_config(
14771593
config_dict: Sequence[Tuple[String, logic.Value]],
14781594
storage_integration_opt: Optional[Sequence[Tuple[String, logic.Value]]],

proto/relationalai/lqp/v1/logic.proto

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

291+
// A single named CSV column with its type. Used to describe both shared key columns and
292+
// per-relation value columns in the generalized `Relations` loading construct.
293+
message NamedColumn {
294+
string name = 1; // CSV column name (e.g. "src"); special name "METADATA$KEY" => derived hash
295+
Type type = 2; // Column type
296+
}
297+
298+
// One output relation: the shared keys plus this relation's own (possibly empty) value columns.
299+
message OutputRelation {
300+
RelationId target_id = 1; // Output relation path
301+
repeated NamedColumn values = 2; // Value columns for this relation (may be empty)
302+
}
303+
304+
// Generalized CSV loading: a shared set of key columns and one or more output relations.
305+
// CDC vs non-CDC is implied by which group is populated:
306+
// - `relations` populated => non-CDC outputs
307+
// - `inserts`/`deletes` => CDC insert/delete groups
308+
message Relations {
309+
repeated NamedColumn keys = 1; // Shared key columns (name "METADATA$KEY" => derived hash)
310+
repeated OutputRelation relations = 2; // Non-CDC outputs
311+
repeated OutputRelation inserts = 3; // CDC insert group
312+
repeated OutputRelation deletes = 4; // CDC delete group
313+
}
314+
291315
message CSVData {
292316
CSVLocator locator = 1;
293317
CSVConfig config = 2;
294318
repeated GNFColumn columns = 3;
295319
string asof = 4; // Blob storage timestamp for freshness requirements
320+
optional Relations relations = 5; // If present, generalized loading; mutually exclusive with columns
296321
}
297322

298323
message CSVLocator {

0 commit comments

Comments
 (0)