The scope rules below have always been what the runtime does. Since Alginte 0.12.0 the editor
enforces two of them it used to let through: an aggregate initializer has no record in
scope, and Map Values sees only the value.
What each expression can see
key and value are the record at the operator’s input. #-variables are extra context some
roles carry. A name that is not in scope is rejected by the editor with a hint, and would fail
deployed.
Two of these bite people:
- Map Values has no
key. The runtime evaluates it against the value alone. To read the key while changing the value, use Map with the key expression set tokey. - The initializer has no record. Kafka calls it once per new key with no arguments, so
value.get('quantity')there has nothing to read. Put the first record’s contribution in the adder, which runs for it like for every other record.
What each serde hands you
The deserializer’s own object reaches your expression, forkey and value alike. There is no
conversion layer, so what the preview shows is what runs.
Strings inside Avro are
org.apache.avro.util.Utf8, not String: == 'Kettle' compares
content and works, .contains('ttl') does not exist on it and fails, .toString() first makes
every String method available. Protobuf and JSON Schema hand you real Strings. An Avro enum
symbol compares through .toString(); a timestamp-millis field is a long and does
arithmetic; a decimal logical type is bytes to the expression and does not compare as a number.
A missing field is loud on Avro and silent on the other two. On Avro and Protobuf an unknown
name is always a typo or schema drift; on JSON Schema an absent key may be an optional property,
which is why the quiet null is the honest answer there. The practical consequence: a filter on
a misspelt field fails visibly on an Avro topic and quietly drops every record on a Protobuf or
JSON Schema one. Whether to make the three agree is an open decision; until it is taken, this
table is the behaviour.
Windowed keys
After a Tumbling, Hopping, Sliding or Session Window and its aggregate,key is
Kafka’s Windowed object all the way through To Stream and beyond, not the plain key. It
has to be unwrapped before a sink, because no key serde here can serialize a Windowed: deployed
as it is, the stream would die on the first record. Alginte refuses that graph, twice: the sink
node grows an error badge while you draw, and the deploy is rejected with the same message, both
naming the fix. Read the key with key.key(), key.window().start() and key.window().end() in
any expression that has the key in scope:
Cookbook
Each recipe shows the expression, the record before and the record after, as its corpus row runs it.orders is an Avro topic (orderId, customerId, item, quantity, priceEur);
shipments is JSON Schema (carrier, weightKg); payments is Protobuf (method).
Filter on a field
Filter, Inclusion Predicate, onorders:
{quantity: 3, ...} passes; {quantity: 1, ...} is dropped. The same predicate is what a
Split branch evaluates.
Reach a nested field
Map Values, on an Avro record whosecustomer is itself a record:
{customer: {id: 'c-1', ...}, ...} becomes {"cid": "c-1"} at a JSON sink. A nested record
placed into the output map, {'c': value.get('customer')}, lands as a JSON object.
Rekey by a field
Select Key, then a downstream filter on the new key:.toString() turns Avro’s Utf8 into a String key; without it the key still matches
'c-6650' and still reaches the sink, but as a Utf8.
Map to a smaller record
Map Values, onorders, keeping three fields and computing one:
{orderId: 'o-1', customerId: 'c-6650', item: 'Kettle', quantity: 3, priceEur: 17.0} becomes
{"customer": "c-6650", "product": "Kettle", "total": 51.0}. Downstream, value.get('total') > 50
passes it and value.get('total') > 100 drops it; the editor knows the computed shape and
completes get('total') for you, marked (inferred).
Several records from one
FlatMap Values, an inline list:"Kettle" and "o-1", with the same key. The word-count
example’s value.toLowerCase().split('\\W+') is the same recipe over a String value.
Count or sum per key
Aggregate overorders, grouped by key. A numeric accumulator:
7 in the store. A map-shaped accumulator, which is what
you want when the aggregate carries more than one figure:
get. A String method on the
accumulator, #aggValue.substring(0, 3), fails deployed on a map-shaped aggregate, and the
editor says so because it evaluates the adder against what your initializer really built.
Concatenate with reduce
Reduce over a String topic:a, b, c for one key reduce to c|b|a. Reduce keeps the value’s type: it cannot change a
String into a map, which is what Aggregate is for.
Join two records
Value Joiner, both sides Avro:Kettle/4. Mixed serdes read each side its own way: #leftValue.get('item') + ':' + #rightValue with a String right side gives Kettle:warehouse-7, and #rightValue.get('item')
on that String right side fails, deployed and in the editor alike, once both join lanes are
wired. Arithmetic between two String sides, #leftValue - #rightValue, fails.
Foreign-key join
Foreign Key Extractor, on the left table’s value, then the joiner:3:Kettle joined to the right table’s Kettle row gives 3:Kettle@warehouse-7. The extractor
sees only value: key there fails.
Unwrap a windowed key
Select Key after a windowed count, as the tumbling-window example does:k1 over seven seconds in 5-second windows produce counts keyed k1@0 and
k1@5000, with the Long count as the value.
Call your own Java
Any expression, once the function is registered (Custom functions):{"masked": "c***@example.com"}. A function that throws fails the expression deployed and
in the editor; an unregistered name fails everywhere too.
Peek without risk
Peek, Action Expression: the result is discarded, and a throw is logged and swallowed, so a peek can never stop the stream. That makes it a place to call a registered function for its side effect,#audit(key, value), and not a place for a guard: an expression that fails in a peek
does not hold the record back.
The forms that do not work
Pinned as failures on every surface, so the editor reports them before you deploy:value['item']on Avro or Protobuf: indexing is aMapoperation and those are not maps. It works on JSON and JSON Schema.value.get('item')on a String topic: a String has no fields.value.get('priceEur').asDouble()on Avro: the field is already a number; there is noasDouble().value.get('amount') > 10on an Avrodecimal: bytes do not compare as a number.#newValuein an adder: the new record isvalue;#newValueis not bound.keyin a Map Values,valuein an initializer,keyin a foreign-key extractor.