defineEtl from the @bijection/pipelines component: the
inputs to capture, the tables to publish, and a transform function. Bijection
reruns the pipeline whenever an input changes and publishes every output table
together.
This page builds the orderSummary pipeline, which publishes the
open_orders and customer_totals tables declared on the
previous page.
Setting up a pipeline
1
Install from npm
Install the pipelines component, and the datasets package whose access
types the policy below uses.
2
Add the pipelines component
Add the component to your app by calling
app.use in the
bijection.config.ts file in your app’s bijection/ folder. The next
bijection dev generates components.pipelines, which the pipeline
definition below refers to. See using components.bijection/bijection.config.ts
3
Decide who may run it
A pipeline checks every operation against an access policy: an internal
query that receives the caller’s Return the same
actor (their tokenIdentifier), the
requested operation and the pipeline as resource. It returns
{ scope } to grant the operation, or null to refuse it.bijection/access.ts
scope for the same actor every time. Background work runs
as the user who configured the pipeline, and the policy is asked again for
that user as the pipeline captures and publishes. If it stops granting
them, the pipeline’s published tables become unreadable.4
Define the pipeline
Call Deploy with
defineEtl and export the functions it generates.bijection/summaries.ts
bijection dev or bijection deploy. Deployment checks that
manifest, steps, page and rebuild name functions your app exports.5
Configure it
Call The first run starts on its own. From then on, every committed change to an
input table starts a new run. Changes that arrive during a run are
collected, and one successor run picks them up. Calling
configure from a mutation run by an authenticated user the policy
grants. It returns the pipeline’s group ID.bijection/summaries.ts
configure again
with an unchanged definition returns the same group; after you change the
definition, call it again to apply the change.The pipeline definition
defineEtl takes these options:
Pipeline, input and output names, table names and key fields must start with a
letter and contain only letters, digits and
_, up to 64 characters.
The publisher declaration
Bijection does not let your code write published rows. Instead, the engine reads a result that is ready to publish through three internal queries, and verifies every row it is handed before writing any of it. A publisher declaration, created withdefinePublisher from bijection/server, names those
queries:
Each function is given as a reference or as a
"module:function" string, and
manifest, steps, page and rebuild must be four different functions. The
declaration must be exported from a module of your app, where deployment
discovers it.
defineEtl builds all of this for you. plan.publisher is the
definePublisher declaration, plan.manifestQuery, plan.stepsQuery and
plan.pageQuery are the three queries, and deploymentMutation() is the
rebuild mutation. The names in the publisher option are where you exported
them.
Only your app’s root may declare a publisher. A component cannot publish into
its own tables this way.
Writing the transform
transform receives an object with one array per input, the carried state
(null on the first page) and { step, final }. Input rows include _id and
_creationTime. It returns { output, next }: one array of rows per output,
and the state for the next page, or null.
Input and output rows are limited to these types:
v.string(),v.int64(),v.number()(finite values only),v.boolean()andv.null()- nested
v.object(...), up to eight levels deep - optional fields with
v.optional(...)
key field
must be a required string, and each key may appear only once in an output over
the whole run. A duplicate key fails the run. An empty output is valid, and
publishes an empty table.
The transform runs in a restricted action:
- It has no network access. A
fetchfails the step. - It has about two seconds of execution time per step.
- It cannot change its input; it works on a copy.
transform throws, the run fails with that error and nothing is published.
The previous publication keeps serving.
Streaming large inputs
Setstream to one input to read it in pages. The transform then runs once per
page and can carry state from one page to the next. Every page comes from the
same snapshot of the input, taken when the run started.
bijection/invoices.ts
meta.final is true on the last page of the input. Each page holds at
most 64 rows of the streamed input, fewer when other inputs are present: every
input that is not streamed is captured in full beside each page and counts
against the same 64-row budget.
The state is a validated value of at most 32 KiB, built from the same types as
rows. It cannot hold arrays or records, so an aggregate over an open-ended set
of keys belongs in the output rows, not in the state.
A streamed input holds at most 8,192 rows, captured within 240 seconds. If the
snapshot expires while pages are being read, the capture starts over from a
fresh snapshot, for at most three attempts in total. A run whose input exceeds 8,192 rows fails
with ETL input exceeds 8192 rows; correcting the input starts a new run.
Checking the complete result
validate declares expectations over the complete result of a run, evaluated
after its last step and before it can be published. A run with a failed
blocking expectation publishes nothing, and the previous publication keeps
serving.
bijection/summaries.ts
reduce(outputs, state, meta)folds the output of each step, in step order, into a validatedstateof at most 32 KiB. The first call receivesnull.check(state, candidate)runs once, after the last fold, and returns a list of verdicts{ expectation, severity, passed, detail? }.candidate.rowsholds each output’s row count.candidate.previousdescribes the serving publication it would replace, or isnull.references, at most two, require everyfromrow’s stringfieldto be a key of thetooutput in the same result.rejects: { member, max }names an output that holds rejected input records. Its table needs a required stringreasonfield. More thanmaxrejected rows blocks the result; any rejected row is an advisory warning. The rejected rows publish with the other outputs, readable under that table’s own access rule.reducenever receives them.
blocking verdict that fails holds the result back with the error
Expectation <name> failed: <detail>. An advisory verdict that fails is kept
as a warning and does not prevent publication. A throw in reduce or check
counts as a failed blocking verdict named validation. You can return at most
32 verdicts, with unique names of 1 to 128 characters that do not contain :,
and details of at most 1 KiB.
Managing runs
TheEtl object provides these operations. Call them from your own
mutations and queries; each one requires an authenticated user and checks the
access policy for them.
Mutations also accept a final
{ request } option. Repeating an operation with
the same request key returns the recorded outcome instead of acting again,
which makes a retry after a lost reply safe.
select returns { status: "requested", run, selection }. Publication happens
afterwards, in its own commit; use publication(ctx) to see when it is
serving. A withdrawn publication stops serving at once, and its tables stay
unavailable until another run publishes. Data a client read earlier cannot be
recalled.
Changing a published group
Changing a pipeline’s transform code needs no special step: the next run uses the new code. Changing which tables the group publishes does. Renaming the pipeline’skey, moving its publisher export to another module or name,
adding, removing or renaming an output, changing an output table’s row type,
or retiring the pipeline all require an explicit transition in the publisher
declaration. A deployment that makes such a change without one is refused and
changes nothing.
bijection/summaries.ts
from names the installed owner exactly: its publisher, producer, complete
list of output tables, and the current owner epoch. Read these from the
deployment’s publisher ownership inspection, a POST to
/api/get_config_hashes on the deployment URL with publisherOwnership: true
and the deployment’s admin key. The CLI does not fill in the epoch or show the
plan before deploying. A transition whose epoch is no longer current is refused
with PublisherOwnership. Set toProducer: null to retire the group; remove
the declaration on a later deployment.
A transition makes the whole group unavailable until the new owner publishes a
complete run. Reads fail with PublishedUnavailable in between; unchanged rows
are recomputed, not carried over. Deployment schedules the rebuild mutation,
which starts that run without any further call. A retired table stays reserved,
unreadable and unwritable, until you delete it.
Limits
Every limit that a definition can exceed is refused when you call
configure,
before it affects running pipelines. The limits on input and output size are
checked as the run proceeds, and a run that exceeds one fails without
publishing.
These limits bound what a run may do; they are not measured capacity.
Throughput at the upper limits has not been measured.