Skip to content
Merged
Show file tree
Hide file tree
Changes from 6 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
124 changes: 120 additions & 4 deletions meta/src/meta/grammar.y
Original file line number Diff line number Diff line change
Expand Up @@ -87,6 +87,14 @@
%nonterm gnf_column logic.GNFColumn
%nonterm gnf_column_path Sequence[String]
%nonterm gnf_columns Sequence[logic.GNFColumn]
%nonterm named_column logic.NamedColumn
%nonterm relation_keys Sequence[logic.NamedColumn]
%nonterm target_relation logic.TargetRelation
%nonterm non_cdc_relations Sequence[logic.TargetRelation]
%nonterm cdc_inserts Sequence[logic.TargetRelation]
%nonterm cdc_deletes Sequence[logic.TargetRelation]
%nonterm relation_body logic.TargetRelations
%nonterm target_relations logic.TargetRelations
%nonterm csv_config logic.CSVConfig
%nonterm csv_data logic.CSVData
%nonterm csv_locator_inline_data String
Expand Down Expand Up @@ -1107,13 +1115,58 @@ csv_asof
: "(" "asof" STRING ")"

csv_data
: "(" "csv_data" csvlocator csv_config gnf_columns csv_asof ")"
construct: $$ = logic.CSVData(locator=$3, config=$4, columns=$5, asof=$6)
: "(" "csv_data" csvlocator csv_config gnf_columns? target_relations? csv_asof ")"
construct: $$ = construct_csv_data($3, $4, $5, $6, $7)
deconstruct:
$3: logic.CSVLocator = $$.locator
$4: logic.CSVConfig = $$.config
$5: Sequence[logic.GNFColumn] = $$.columns
$6: String = $$.asof
$5: Optional[Sequence[logic.GNFColumn]] = deconstruct_csv_data_columns_optional($$)
$6: Optional[logic.TargetRelations] = deconstruct_csv_data_relations_optional($$)
$7: String = $$.asof

named_column
: "(" "column" STRING type ")"
construct: $$ = logic.NamedColumn(name=$3, type=$4)
deconstruct:
$3: String = $$.name
$4: logic.Type = $$.type

relation_keys
: "(" "keys" named_column* ")"

target_relation
: "(" "relation" relation_id named_column* ")"
construct: $$ = logic.TargetRelation(target_id=$3, values=$4)
deconstruct:
$3: logic.RelationId = $$.target_id
$4: Sequence[logic.NamedColumn] = $$.values

non_cdc_relations
: target_relation*

cdc_inserts
: "(" "inserts" target_relation* ")"

cdc_deletes
: "(" "deletes" target_relation* ")"

relation_body
: non_cdc_relations
construct: $$ = construct_non_cdc_relations($1)
deconstruct if builtin.is_empty($$.inserts) and builtin.is_empty($$.deletes):
$1: Sequence[logic.TargetRelation] = $$.relations
| cdc_inserts cdc_deletes
construct: $$ = construct_cdc_relations($1, $2)
deconstruct if not (builtin.is_empty($$.inserts) and builtin.is_empty($$.deletes)):
$1: Sequence[logic.TargetRelation] = $$.inserts
$2: Sequence[logic.TargetRelation] = $$.deletes

target_relations
: "(" "relations" relation_keys relation_body ")"
construct: $$ = construct_relations($3, $4)
deconstruct:
$3: Sequence[logic.NamedColumn] = $$.keys
$4: logic.TargetRelations = $$

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


def construct_non_cdc_relations(
relations: Sequence[logic.TargetRelation],
) -> logic.TargetRelations:
return logic.TargetRelations(
keys=list[logic.NamedColumn](),
relations=relations,
inserts=list[logic.TargetRelation](),
deletes=list[logic.TargetRelation](),
)


def construct_cdc_relations(
inserts: Sequence[logic.TargetRelation],
deletes: Sequence[logic.TargetRelation],
) -> logic.TargetRelations:
return logic.TargetRelations(
keys=list[logic.NamedColumn](),
relations=list[logic.TargetRelation](),
inserts=inserts,
deletes=deletes,
)


def construct_relations(
keys: Sequence[logic.NamedColumn],
body: logic.TargetRelations,
) -> logic.TargetRelations:
return logic.TargetRelations(
keys=keys,
relations=body.relations,
inserts=body.inserts,
deletes=body.deletes,
)


def construct_csv_data(
locator: logic.CSVLocator,
config: logic.CSVConfig,
columns_opt: Optional[Sequence[logic.GNFColumn]],
relations_opt: Optional[logic.TargetRelations],
asof: String,
) -> logic.CSVData:
return logic.CSVData(
locator=locator,
config=config,
columns=builtin.unwrap_option_or(columns_opt, list[logic.GNFColumn]()),
asof=asof,
relations=relations_opt,
)


def deconstruct_csv_data_columns_optional(msg: logic.CSVData) -> Optional[Sequence[logic.GNFColumn]]:
if builtin.has_proto_field(msg, "relations"):
return builtin.none()
return builtin.some(msg.columns)


def deconstruct_csv_data_relations_optional(msg: logic.CSVData) -> Optional[logic.TargetRelations]:
if builtin.has_proto_field(msg, "relations"):
return builtin.some(builtin.unwrap_option(msg.relations))
return builtin.none()


def construct_csv_config(
config_dict: Sequence[Tuple[String, logic.Value]],
storage_integration_opt: Optional[Sequence[Tuple[String, logic.Value]]],
Expand Down
25 changes: 25 additions & 0 deletions proto/relationalai/lqp/v1/logic.proto
Original file line number Diff line number Diff line change
Expand Up @@ -288,11 +288,36 @@ message StorageIntegration {
string s3_secret_access_key = 5;
}

// A single named column with its type. Used to describe both shared key columns and
// per-relation value columns in the generalized `TargetRelations` loading construct.
message NamedColumn {
string name = 1; // Column name (e.g. "src"); special name "METADATA$KEY" => derived hash
Comment thread
hbarthels marked this conversation as resolved.
Outdated
Type type = 2; // Column type
}

// One target relation: the shared keys plus this relation's own (possibly empty) value columns.
message TargetRelation {
RelationId target_id = 1; // Output relation path
repeated NamedColumn values = 2; // Value columns for this relation (may be empty)
}

// Generalized loading: a shared set of key columns and one or more target relations.
// CDC vs non-CDC is implied by which group is populated:
// - `relations` populated => non-CDC outputs
// - `inserts`/`deletes` => CDC insert/delete groups

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Does this mean that the relations and the inserts/deletes portion of TargetRelations are mutually exclusive? e.g. either one or the other is populated? If so, maybe we can express this in a OneOf of CDC/non-CDC?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yeah, they should be mutually exclusive. I will look into that.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I changed it to a OneOf.

message TargetRelations {
repeated NamedColumn keys = 1; // Shared key columns (name "METADATA$KEY" => derived hash)
repeated TargetRelation relations = 2; // Non-CDC outputs
repeated TargetRelation inserts = 3; // CDC insert group
repeated TargetRelation deletes = 4; // CDC delete group
}

message CSVData {
CSVLocator locator = 1;
CSVConfig config = 2;
repeated GNFColumn columns = 3;
string asof = 4; // Blob storage timestamp for freshness requirements
optional TargetRelations relations = 5; // If present, generalized loading; mutually exclusive with columns
}

message CSVLocator {
Expand Down
Loading
Loading