How do I combine manifests when one source decides the type per row?¶
This is the plant of the manifest union example (20), with one difference. The
maintenance system keeps a single register for everything it maintains, and a
type column says whether a row is a machine or a production line. Its
manifest sends each row to its type with a router. The sensor feed reports
devices, as in that example, and a device is the same machine as a register
row when their serial numbers agree.
You want machines and devices to become one type, Machine, matched on the
serial number, while production lines keep going to their own type. Splitting
the register into one resource per type would read the table twice and repeat
what the type column already says. When GraFlo combines the manifests, the
router stays whole, and the matching applies only to the rows the router sends
to Machine.
flowchart LR
register[register.csv] --> router{type column}
router -- machine --> Machine
router -- line --> ProductionLine
devices[devices.csv] -- same serial number --> Machine
What you need¶
- GraFlo installed (
pip install graflo). No database is needed. - The merge declarations of Combine two manifests and the router of One table that holds many kinds of things.
The data¶
data/register.csv, from the maintenance system:
| type | asset_id | serial_number | name |
|---|---|---|---|
| machine | A1 | HP-0042 | Hydraulic press |
| machine | A2 | Conveyor | |
| line | L1 | Assembly line 1 |
data/devices.csv is the sensor feed of example 20: D7
(hp-0042) is the hydraulic press, and D9 (LT-0007) is a lathe the
register does not list.
Steps¶
1. Route the register rows by type¶
The register resource of manifest_maintenance.yaml
has one step. The router reads the type column of each row and makes the row
a vertex of the type that type_map names:
- name: register
pipeline:
- vertex_router:
type_field: type
type_map:
machine: Machine
line: ProductionLine
The sensor feed's manifest, manifest_sensors.yaml,
is the one from example 20.
2. Declare how the manifests combine¶
merge.yaml holds the same three declarations as example 20.
The maintenance manifest already calls its type Machine, so the
equivalence names the combined type with into instead of a map of names:
canonical_maps:
right: {properties: {Device: {serial: serial_number}}}
vertex_equivalences:
- left: Machine
right: Device
into: Machine
identity:
- name: match_key
sources:
register: {input: [serial_number]}
devices: {input: [serial]}
- local_key:
register: {field: asset_id, tag: maintenance}
devices: {field: device_id, tag: sensors}
The identity names the register resource as a whole. Nothing in it
mentions the type column or production lines.
3. Build the combined manifest¶
cd examples/21-router-union-alignment
uv run graflo merge manifest_maintenance.yaml manifest_sensors.yaml \
--op merge.yaml -o artifacts/manifest_union.yaml
In artifacts/manifest_union.yaml the
register resource keeps its one router, followed by the steps that compute
the keys:
- vertex_router:
type_field: type
type_map:
machine: Machine
line: ProductionLine
Machine: Machine
ProductionLine: ProductionLine
type_map_only: true
- transform:
call:
module: graflo.util.transform
foo: normalized_key
params: {}
output:
- match_key
input:
- serial_number
when:
field: type
in:
- machine
# the step that computes local_key follows, with the same `when`
The router is not split. GraFlo adds the type names themselves to type_map
and closes the table with type_map_only: true, so the router accepts the same
type values as before and no others.
Each key step carries when: {field: type, in: [machine]}. GraFlo derived that
condition from the router: the step runs only for the rows the router sends to
Machine, and a production line row never gets a key.
What you should see¶
uv run python inspect_fusion.py runs both files through the combined manifest
and prints the vertex each record becomes:
resource own key vertex type matched on vertex id
register A1 Machine hp-0042 303d50890862
register A2 Machine maintenance:A2 d154517d907c
register L1 ProductionLine - L1
devices D7 Machine hp-0042 303d50890862
devices D9 Machine lt-0007 27ded24b71df
4 machine records -> 3 vertices
The machine row A1 is matched on its serial number and lands on the same
vertex as the device D7. A2 has no serial number and falls back to its own
key. The line row L1 goes through the same router to ProductionLine and
keeps its own identity, asset_id. The three machine vertices are the ones
example 20 produces.
Also possible¶
If the router sent two of its types into the combined type, say machine rows
to Machine and robot rows to Robot, the equivalence would list both
(left: [Machine, Robot]). The register's entry in sources can then be keyed
by type, and each type gets its own step, guarded on its own type value:
A shared vocabulary may already call several register types by one name.
merge_vocabulary.yaml calls machines and production
lines Equipment, and keeps the equivalence about machines and devices:
canonical_maps:
left:
vertices: {Machine: Equipment, ProductionLine: Equipment}
allow_merges: true
vertex_equivalences:
- left: Machine
right: Device
The combined type is Equipment, and production lines are part of it without
being listed. Only machines are matched with devices: the derivations are keyed
by Machine, and a production line keeps its own key behind a tag:
naming (vertex):
merged left right via
Equipment Machine Device equivalence
Equipment ProductionLine - vocabulary
resource own key vertex type matched on vertex id
register A1 Equipment hp-0042 303d50890862
register A2 Equipment maintenance:A2 d154517d907c
register L1 Equipment left:ProductionLine:L1 cb413d9d4729
devices D7 Equipment hp-0042 303d50890862
devices D9 Equipment lt-0007 27ded24b71df
See a vocabulary that merges several types.
What to read next¶
- Version control for a manifest: two people changed the same manifest, and their changes are merged against the version they started from.
- Routed sources: why the union closes a router over its own side's types.
Files¶
The example lives in examples/21-router-union-alignment.
manifest_maintenance.yaml
schema:
metadata:
name: maintenance
version: 1.0.0
graph:
vertex_config:
vertices:
- name: Machine
properties:
- asset_id
- serial_number
- name
identity:
- asset_id
- name: ProductionLine
properties:
- asset_id
- name
identity:
- asset_id
edge_config:
edges: []
ingestion_model:
resources:
# One register holds machines and production lines; its `type` column says
# which type each row becomes.
- name: register
pipeline:
- vertex_router:
type_field: type
type_map:
machine: Machine
line: ProductionLine
transforms: []
manifest_sensors.yaml
merge.yaml
# How to combine manifest_maintenance.yaml (left) and manifest_sensors.yaml (right).
op: merge_manifests
# 1. A device's `serial` fills a machine's `serial_number`.
canonical_maps:
right: {properties: {Device: {serial: serial_number}}}
# 2. Machines and devices are one type, named Machine, and this is how records
# from both sides find each other.
vertex_equivalences:
- left: Machine
right: Device
into: Machine
identity:
- name: match_key
sources:
register: {input: [serial_number]}
devices: {input: [serial]}
- local_key:
register: {field: asset_id, tag: maintenance}
devices: {field: device_id, tag: sensors}
merge_vocabulary.yaml
# A shared vocabulary calls machines and production lines Equipment; only
# machines are matched with devices.
op: merge_manifests
# 1. The vocabulary: both register types are Equipment, and a device's
# `serial` is a `serial_number`.
canonical_maps:
left:
vertices: {Machine: Equipment, ProductionLine: Equipment}
allow_merges: true
right: {properties: {Device: {serial: serial_number}}}
# 2. Machines and devices are one type. The vocabulary already makes production
# lines Equipment too, so they join it without being listed here, and `into`
# comes from the vocabulary.
vertex_equivalences:
- left: Machine
right: Device
identity:
- name: match_key
sources:
register: {Machine: {input: [serial_number]}}
devices: {input: [serial]}
- local_key:
register: {Machine: {field: asset_id, tag: maintenance}}
devices: {field: device_id, tag: sensors}
inspect_fusion.py
"""Show where each register row and each device lands in the combined manifest.
Reads the combined manifest, runs the two CSV files through it and prints, for
every record, the vertex type it becomes, the matching key derived for it, and
the vertex id. No database is involved.
cd examples/21-router-union-alignment
uv run python inspect_fusion.py
"""
from __future__ import annotations
import asyncio
import csv
from pathlib import Path
from typing import Any
from suthing import FileHandle
from graflo import GraphManifest
from graflo.hq.document_caster import DocumentCaster
from graflo.hq.ingestion_parameters import IngestionParams
EXAMPLE_DIR = Path(__file__).resolve().parent
#: Resource name, its file, and the column holding the record's own key.
SOURCES = [
("register", "register.csv", "asset_id"),
("devices", "devices.csv", "device_id"),
]
ID_WIDTH = 12
def read_rows(name: str) -> list[dict[str, str]]:
"""Read one CSV file of the example as a list of rows."""
with open(EXAMPLE_DIR / "data" / name, newline="") as f:
return list(csv.DictReader(f))
def cast(caster: DocumentCaster, resource: str, filename: str) -> Any:
"""Run one file through one resource and return the resulting graph."""
rows = read_rows(filename)
result = asyncio.run(caster.cast_batch(rows, resource, params=IngestionParams()))
return result.graph
def main() -> None:
"""Print every vertex produced, then the machine count."""
manifest = GraphManifest.from_config(
FileHandle.load(EXAMPLE_DIR / "artifacts" / "manifest_union.yaml")
)
manifest.finish_init()
caster = DocumentCaster(manifest.require_ingestion_model())
machine_ids: list[str] = []
print(
f"{'resource':<10}{'own key':<9}{'vertex type':<16}{'matched on':<16}vertex id"
)
for resource, filename, key in SOURCES:
for vertex_type, docs in cast(caster, resource, filename).vertices.items():
for doc in docs:
# The derived key a machine is matched on. A production line has
# none: it keeps its own identity, asset_id.
matched_on = doc.get("match_key") or doc.get("local_key") or "-"
# A machine is keyed by a digest; a production line by its asset_id.
vertex_id = doc.get("id", doc[key])[:ID_WIDTH]
print(
f"{resource:<10}{doc[key]:<9}{vertex_type:<16}{matched_on:<16}"
f"{vertex_id}"
)
if vertex_type == "Machine":
machine_ids.append(vertex_id)
print(f"{len(machine_ids)} machine records -> {len(set(machine_ids))} vertices")
if __name__ == "__main__":
main()