Skip to content

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

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
schema: 2 vertices, 0 edges, version 1.1.0
resources: 2
written: 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:

register:
    Machine: {input: [serial_number]}
    Robot: {input: [serial_number]}

See one derivation per type.

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.

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
schema:
    metadata:
        name: sensors
        version: 1.0.0
    graph:
        vertex_config:
            vertices:
            -   name: Device
                properties:
                -   device_id
                -   serial
                -   model
                identity:
                -   device_id
        edge_config:
            edges: []
ingestion_model:
    resources:
    -   name: devices
        pipeline:
        -   vertex: Device
    transforms: []
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()