Standalone, Discord-free Java library for authoring typed data-extraction pipelines over arbitrary HTML / XML / JSON / Lua / text input.
Pipelines are flat sequences of stages - Source -> (Filter | Transform | Predicate)* -> Terminal -
with sub-pipeline bodies (per-element Map, named fan-out through COLLECT_MAP and
TRANSFORM_JSON_OBJECT_BUILD), pipeline operands that read a second document beside the running
value, and EmbedSource (run a saved pipeline as a single source-like stage). Modeled after
Java 8 Streams; the stage taxonomy mirrors Stream operation categories.
// build.gradle.kts
dependencies {
implementation("com.github.simplified-dev:dataflow:master-SNAPSHOT")
}Requires Java 21.
import dev.simplified.dataflow.DataPipeline;
import dev.simplified.dataflow.DataTypes;
import dev.simplified.dataflow.PipelineContext;
import dev.simplified.dataflow.stage.filter.dom.TextContainsFilter;
import dev.simplified.dataflow.stage.source.UrlSource;
import dev.simplified.dataflow.stage.terminal.collect.FirstCollect;
import dev.simplified.dataflow.stage.transform.dom.CssSelectTransform;
import dev.simplified.dataflow.stage.transform.dom.NthChildTransform;
import dev.simplified.dataflow.stage.transform.dom.ParseHtmlTransform;
import dev.simplified.dataflow.stage.transform.dom.TextTransform;
import dev.simplified.dataflow.stage.transform.primitive.ParseIntTransform;
import dev.simplified.dataflow.stage.transform.string.RegexExtractTransform;
DataPipeline<Integer> pipeline = DataPipeline.builder()
.source(UrlSource.rawHtml("https://hypixelskyblock.minecraft.wiki/w/Dark_Claymore"))
.stage(ParseHtmlTransform.of())
.stage(CssSelectTransform.of("table.infobox tr"))
.stage(TextContainsFilter.of("Dmg"))
.stage(FirstCollect.of(DataTypes.DOM_NODE))
.stage(NthChildTransform.of("td", 1))
.stage(TextTransform.of())
.stage(RegexExtractTransform.of("\\d+"))
.stage(ParseIntTransform.of())
.build();
Integer dmg = pipeline.execute(PipelineContext.defaults()); // 500Class names drop their role suffix here (Contains is ContainsFilter under .filter and
ContainsPredicate under .predicate).
dev.simplified.dataflow.stage
.source Url, Literal, LiteralList, Embed
.filter.string Contains, Matches, StartsWith, EndsWith, Equals, NonEmpty
.filter.list Distinct, DistinctBy, NotNull, Take, Skip, IndexInRange,
TakeWhile, DropWhile, Where
.filter.numeric Int/Long/Double x {GreaterThan, LessThan, InRange}
.filter.dom TextContains, TextMatches, HasAttr, TagEquals
.filter.json HasField, FieldEquals
.transform.string LowerCase, UpperCase, CamelCase, PascalCase, SnakeCase, TitleCase,
Length, Prepend, Append, Trim, Replace, ReplaceMatch, Split,
RegexExtract, ValueMap, RangeExpand, Fetch
.transform.primitive ParseInt/Long/Float/Double/Boolean/Roman, Abs*, Negate*,
Arithmetic*, BinaryArithmetic*, RoundFloat/Double, Constant,
Coalesce, ToString, ToRaw, Peek, Expect (+ ArithmeticOperator)
.transform.list Size, Reverse, Sort, SortBy, Map, FlatMap, Flatten, Concat,
GroupBy, Enumerate, Zip, Rotate, Broadcast
.transform.dom ParseHtml, CssSelect, Text, OwnText, Attr, NthChild, Children,
Parent, OuterHtml, SpanExpand
.transform.json ParseJson, ParseXml, ParseLua, Path, Field,
AsString/Int/Long/Double/Boolean, Stringify, Deserialize,
ObjectBuild, Entries, JoinByKey, KeyLookup, ResolveAncestor
.transform.encoding Base64Encode/Decode, UrlEncode/Decode, HtmlDecode, JsonUnescape
.predicate.string Contains, Matches, StartsWith, EndsWith, Equals, NonEmpty
.predicate.numeric Int/Long/Double x {GreaterThan, LessThan, InRange}
.predicate.dom TextContains, TextMatches, HasAttr, TagEquals
.predicate.json HasField, FieldEquals
.predicate.common NotNull, Not, And, Or, Compare (+ CompareOperator)
.terminal.collect First, Last, Nth, List, SubList, Set, Join, Map,
JsonObjectFromEntries
.terminal.sum Count, SumInt, SumLong, SumDouble
.terminal.average AverageInt, AverageLong, AverageDouble
.terminal.minmax Min, Max, MinBy, MaxBy
.terminal.match AnyMatch, AllMatch, NoneMatch, FindFirst
DataTypes: NONE, RAW_HTML, RAW_XML, RAW_JSON, STRING, INT, LONG,
FLOAT, DOUBLE, BOOLEAN, DOM_NODE, JSON_ELEMENT, JSON_OBJECT, JSON_ARRAY,
MAP_OUTPUT, plus List<X> / Set<X> constructors and any type a host adds with
DataTypes.register.
Every kind below is the @StageSpec(id = ...) of one class, resolved on load by
StageRegistry.byId. The join, lookup, ancestor, group, enumerate, zip and broadcast stages
never mutate the rows they are given: what they emit is copied.
| Kind | Input | Output | Notes |
|---|---|---|---|
SOURCE_URL |
NONE | RAW_* / STRING |
ctx.fetcher(); optional maxBodyBytes; any status outside 2xx (an unfollowed 3xx, a 4xx, a 5xx), or a ctx.fetchGuard() refusal, fails the run, except that a 5xx to a refresh of a stale cached copy inside its stale-if-error window answers the cached body (see Fetching) |
SOURCE_LITERAL |
NONE | T |
literal value parsed from config |
SOURCE_LITERAL_LIST |
NONE | List<T> |
literal list from JSON array |
SOURCE_EMBED |
NONE | declared at construction | resolves saved pipeline by id |
| Kind | Input | Output | Notes |
|---|---|---|---|
PARSE_HTML |
RAW_HTML |
DOM_NODE |
|
PARSE_XML |
RAW_XML |
JSON_ELEMENT |
|
PARSE_JSON |
RAW_JSON |
JSON_ELEMENT |
|
PARSE_LUA |
STRING |
JSON_ELEMENT |
a Lua 5.1 data module: locals, then return of a literal table; anything computed throws with its line and column |
TRANSFORM_CSS_SELECT |
DOM_NODE |
List<DOM_NODE> |
|
TRANSFORM_DOM_TEXT |
DOM_NODE |
STRING |
|
TRANSFORM_DOM_ATTR |
DOM_NODE |
STRING |
|
TRANSFORM_DOM_NTH_CHILD |
DOM_NODE |
DOM_NODE |
|
TRANSFORM_DOM_CHILDREN |
DOM_NODE |
List<DOM_NODE> |
|
TRANSFORM_DOM_PARENT |
DOM_NODE |
DOM_NODE |
|
TRANSFORM_DOM_OUTER_HTML |
DOM_NODE |
STRING |
|
TRANSFORM_DOM_OWN_TEXT |
DOM_NODE |
STRING |
|
TRANSFORM_DOM_SPAN_EXPAND |
DOM_NODE |
List<List<DOM_NODE>> |
a table as a grid, each cell at every position its rowspan / colspan covers, and an empty detached td at a position before a row's last covered one that no cell covers; optional rowSelector picks the rows returned, every row of their row groups still laid out so a span covers the rows HTML shows it covering; a non-table, or a layout past 1,000,000 positions (MAX_ENTRIES, counting every row laid out), rejects |
TRANSFORM_JSON_PATH |
JSON_ELEMENT |
JSON_ELEMENT |
|
TRANSFORM_JSON_FIELD |
JSON_OBJECT |
JSON_ELEMENT |
|
TRANSFORM_JSON_AS_STRING |
JSON_ELEMENT |
STRING |
|
TRANSFORM_JSON_AS_INT |
JSON_ELEMENT |
INT |
|
TRANSFORM_JSON_AS_LONG |
JSON_ELEMENT |
LONG |
|
TRANSFORM_JSON_AS_DOUBLE |
JSON_ELEMENT |
DOUBLE |
|
TRANSFORM_JSON_AS_BOOLEAN |
JSON_ELEMENT |
BOOLEAN |
|
TRANSFORM_JSON_STRINGIFY |
JSON_ELEMENT |
STRING |
compact JSON, no HTML escaping; an object member holding JSON null is written as null, so the text parses back to the input |
TRANSFORM_JSON_DESERIALIZE |
JSON_* |
Basic T or List<X> |
Gson into a Basic type; for List<X> (Basic X) a JSON array becomes a list in array order, a null or unconvertible element is dropped, and a non-array rejects; List<DOM_NODE>, List<NONE> and a List output over a JSON_OBJECT input fail the load |
TRANSFORM_JSON_OBJECT_BUILD |
I |
JSON_OBJECT |
one field per typed named body, in declared order; a null result omits its field |
TRANSFORM_JSON_ENTRIES |
JSON_OBJECT |
List<JSON_OBJECT> |
one {key, value} object per entry, in document order; inline spreads an object value's fields; optional keyField / valueField |
TRANSFORM_JOIN_BY_KEY |
List<JSON_OBJECT> |
List<JSON_OBJECT> |
merges a second document's rows by key (PIPELINE operand right); INNER / LEFT / FULL; a right value fills only an empty left cell, never the left key; FULL appends the first right row of each unmatched key, never a repeated or keyless one; optional columns |
TRANSFORM_KEY_LOOKUP |
STRING |
JSON_ELEMENT |
valueField of the row whose keyField is the input, in a second document (PIPELINE operand table); a miss rejects |
TRANSFORM_RESOLVE_ANCESTOR |
List<JSON_OBJECT> |
List<JSON_OBJECT> |
follows each row's parentField to its root and writes the root's key or valueField into outputField; optional depthField; a cycle writes neither |
TRANSFORM_BASE64_ENCODE |
STRING |
STRING |
|
TRANSFORM_BASE64_DECODE |
STRING |
STRING |
|
TRANSFORM_URL_ENCODE |
STRING |
STRING |
UTF-8 form encoding: a space becomes +, which only a query value reads as a space, so it encodes a query value and not a path segment |
TRANSFORM_URL_DECODE |
STRING |
STRING |
|
TRANSFORM_HTML_DECODE |
STRING |
STRING |
HTML character references, decoded once, as jsoup's Parser.unescapeEntities decodes them; � gives U+0000 and a surrogate reference a lone surrogate, where a browser gives U+FFFD |
TRANSFORM_JSON_UNESCAPE |
STRING |
STRING |
JSON string escapes, surrogate pairs joined; a malformed escape rejects |
TRANSFORM_JOIN_BY_KEY, TRANSFORM_KEY_LOOKUP and TRANSFORM_RESOLVE_ANCESTOR match keys by
their string form - a primitive's text, the compact JSON of an array or object - so 1 and "1"
are one key, 1 and 1.0 are two, and object keys match only with their members in the same
order. TRANSFORM_GROUP_BY compares keys as JSON values instead - a number by its exact value
whatever its size, so 1, 1.0 and 1e0 are one key and "1" another, and an object whatever
its member order - and a folded row keeps the key its group's first row wrote, so a join on that
key afterwards matches 1 or 1.0 depending on which row came first.
| Kind | Input | Output | Notes |
|---|---|---|---|
TRANSFORM_REGEX_EXTRACT |
STRING |
STRING |
|
TRANSFORM_TRIM |
STRING |
STRING |
|
TRANSFORM_REPLACE |
STRING |
STRING |
regex, fixed replacement |
TRANSFORM_REPLACE_MATCH |
STRING |
STRING |
each match's capture group rewritten through a body (STRING -> STRING), inserted literally; a null result keeps the match |
TRANSFORM_SPLIT |
STRING |
List<STRING> |
|
TRANSFORM_LOWERCASE |
STRING |
STRING |
|
TRANSFORM_UPPERCASE |
STRING |
STRING |
|
TRANSFORM_CAMEL_CASE |
STRING |
STRING |
|
TRANSFORM_PASCAL_CASE |
STRING |
STRING |
|
TRANSFORM_SNAKE_CASE |
STRING |
STRING |
|
TRANSFORM_TITLE_CASE |
STRING |
STRING |
|
TRANSFORM_STRING_LENGTH |
STRING |
INT |
|
TRANSFORM_PREFIX |
STRING |
STRING |
|
TRANSFORM_SUFFIX |
STRING |
STRING |
|
TRANSFORM_VALUE_MAP |
STRING |
STRING |
exact lookup in a STRING_MAP table; a miss takes defaultValue, rejects under strict, or passes through |
TRANSFORM_RANGE_EXPAND |
STRING |
List<INT> |
"1-15" to [1, ..., 15]; optional regex, step, maxSize (1000; at most 100000, a larger one fails the load) |
TRANSFORM_FETCH |
STRING |
RAW_* / STRING |
fetches the URL the input names, or urlTemplate with {} replaced by it; a ClientError other than 408 and 429 rejects, any other failure throws; optional maxBodyBytes (see Fetching) |
| Kind | Input | Output | Notes |
|---|---|---|---|
TRANSFORM_LIST_LENGTH |
List<T> |
INT |
|
TRANSFORM_REVERSE |
List<T> |
List<T> |
|
TRANSFORM_SORT |
List<T> |
List<T> |
natural ordering, asc/desc |
TRANSFORM_SORT_BY |
List<T> |
List<T> |
key-extractor body (T -> K) |
TRANSFORM_MAP |
List<X> |
List<Y> |
body X -> Y per element; a null result drops the element |
TRANSFORM_FLAT_MAP |
List<X> |
List<Y> |
body X -> List<Y>, concatenated |
TRANSFORM_FLATTEN |
List<List<T>> |
List<T> |
|
TRANSFORM_CONCAT |
List<T> |
List<T> |
appends a second document's elements (PIPELINE operand other) |
TRANSFORM_GROUP_BY |
List<JSON_OBJECT> |
List<JSON_OBJECT> |
one row per keyField value, first-seen order; per-field aggregates (STRING_MAP of FIRST, LAST, LIST, UNION, CONCAT, MAX, MIN, COUNT, MODE), else first non-null |
TRANSFORM_ENUMERATE |
List<T> |
List<JSON_OBJECT> |
{index, value} per element; optional start, indexKey, valueKey; a null, JSON null or NaN / infinite FLOAT / DOUBLE element is dropped but keeps its index |
TRANSFORM_ZIP |
I |
List<JSON_OBJECT> |
pairs two lists read out of one value (bodies left, right) by position; SHORTEST / LONGEST; a null or NaN / infinite FLOAT / DOUBLE element omits its key |
TRANSFORM_ROTATE |
I |
List<T> |
a list read out of the value (body list) rotated by an offset read out of it (body offset) |
TRANSFORM_BROADCAST |
I |
List<JSON_OBJECT> |
one parent value (body parent) carried onto every child of a list (body children); a null or NaN / infinite FLOAT / DOUBLE parent omits its key, and such a child is dropped |
| Kind | Input | Output | Notes |
|---|---|---|---|
TRANSFORM_PARSE_INT |
STRING |
INT |
|
TRANSFORM_PARSE_LONG |
STRING |
LONG |
|
TRANSFORM_PARSE_FLOAT |
STRING |
FLOAT |
|
TRANSFORM_PARSE_DOUBLE |
STRING |
DOUBLE |
|
TRANSFORM_PARSE_BOOLEAN |
STRING |
BOOLEAN |
|
TRANSFORM_PARSE_ROMAN |
STRING |
INT |
canonical numerals I to MMMCMXCIX, either case; anything else rejects |
TRANSFORM_ABS_INT |
INT |
INT |
|
TRANSFORM_ABS_LONG |
LONG |
LONG |
|
TRANSFORM_ABS_FLOAT |
FLOAT |
FLOAT |
|
TRANSFORM_ABS_DOUBLE |
DOUBLE |
DOUBLE |
|
TRANSFORM_NEGATE_INT |
INT |
INT |
|
TRANSFORM_NEGATE_LONG |
LONG |
LONG |
|
TRANSFORM_NEGATE_FLOAT |
FLOAT |
FLOAT |
|
TRANSFORM_NEGATE_DOUBLE |
DOUBLE |
DOUBLE |
|
TRANSFORM_ARITHMETIC_INT |
INT |
INT |
input OP operand, operand an INT |
TRANSFORM_ARITHMETIC_LONG |
LONG |
LONG |
input OP operand, operand a LONG |
TRANSFORM_ARITHMETIC_FLOAT |
FLOAT |
FLOAT |
input OP operand, operand a DOUBLE narrowed to float |
TRANSFORM_ARITHMETIC_DOUBLE |
DOUBLE |
DOUBLE |
input OP operand, operand a DOUBLE |
TRANSFORM_BINARY_ARITHMETIC_INT |
I |
INT |
left OP right, both bodies I -> INT |
TRANSFORM_BINARY_ARITHMETIC_LONG |
I |
LONG |
left OP right, both bodies I -> LONG |
TRANSFORM_BINARY_ARITHMETIC_FLOAT |
I |
FLOAT |
left OP right, both bodies I -> FLOAT |
TRANSFORM_BINARY_ARITHMETIC_DOUBLE |
I |
DOUBLE |
left OP right, both bodies I -> DOUBLE |
TRANSFORM_ROUND_FLOAT |
FLOAT |
FLOAT |
scale decimal places, optional mode (RoundingMode, HALF_UP) |
TRANSFORM_ROUND_DOUBLE |
DOUBLE |
DOUBLE |
scale decimal places, optional mode (RoundingMode, HALF_UP) |
TRANSFORM_CONSTANT |
I |
T |
every present input replaced by value, parsed under the types SOURCE_LITERAL admits but more strictly: a BOOLEAN is true or false in any case, and a FLOAT or DOUBLE must be finite, so a BOOLEAN of yes or an empty string, or a DOUBLE of NaN or 1e309, fails the load |
TRANSFORM_COALESCE |
I |
O |
first non-null of body, fallback body and defaultValue; optional when body guards body; defaultValue is parsed as TRANSFORM_CONSTANT parses value |
TRANSFORM_ITERATE |
T |
T |
body applied to its own output while the optional while body yields true, or until a pass changes nothing; maxIterations (1 to 10,000) still changing fails the run, as does a pass yielding over 10,000,000 characters or elements |
TRANSFORM_TO_STRING |
T |
STRING |
|
TRANSFORM_TO_RAW |
STRING |
RAW_* |
retypes the string; parses nothing |
TRANSFORM_PEEK |
T |
T |
identity + ctx.log() side effect |
TRANSFORM_EXPECT |
T |
T |
identity while the predicate body (T -> BOOLEAN) yields true; false or null fails the run with ExpectationFailedException naming the expectation sentence (see Expectations) |
OP is an ArithmeticOperator named by its constant: ADD, SUBTRACT, MULTIPLY, DIVIDE,
MODULO. INT and LONG compute through Math.*Exact, so an overflow throws
ArithmeticException, and DIVIDE truncates toward zero. MODULO is a floored remainder
(Math.floorMod) for every type: zero or of the divisor's sign, and smaller in magnitude than the
divisor. A FLOAT or DOUBLE remainder that would round onto the divisor (-1e-14 MODULO 360)
is zero with the divisor's sign. A zero divisor under DIVIDE or MODULO, or a NaN or
infinite FLOAT / DOUBLE result, rejects with null. A binary stage rejects when either body
yields null, and the right body does not run once the left one has yielded null. A fraction
in an INT or LONG operand, or a FLOAT operand that is not finite or narrows to zero, fails
the load.
| Kind | Input | Output | Notes |
|---|---|---|---|
FILTER_DOM_TEXT_CONTAINS |
List<DOM_NODE> |
List<DOM_NODE> |
|
FILTER_DOM_TEXT_MATCHES |
List<DOM_NODE> |
List<DOM_NODE> |
|
FILTER_DOM_HAS_ATTR |
List<DOM_NODE> |
List<DOM_NODE> |
|
FILTER_DOM_TAG_EQUALS |
List<DOM_NODE> |
List<DOM_NODE> |
|
FILTER_STRING_CONTAINS |
List<STRING> |
List<STRING> |
|
FILTER_STRING_MATCHES |
List<STRING> |
List<STRING> |
|
FILTER_STRING_STARTS_WITH |
List<STRING> |
List<STRING> |
|
FILTER_STRING_ENDS_WITH |
List<STRING> |
List<STRING> |
|
FILTER_STRING_EQUALS |
List<STRING> |
List<STRING> |
|
FILTER_STRING_NON_EMPTY |
List<STRING> |
List<STRING> |
|
FILTER_JSON_HAS_FIELD |
List<JSON_OBJECT> |
List<JSON_OBJECT> |
|
FILTER_JSON_FIELD_EQUALS |
List<JSON_OBJECT> |
List<JSON_OBJECT> |
|
FILTER_INT_GREATER_THAN |
List<INT> |
List<INT> |
|
FILTER_INT_LESS_THAN |
List<INT> |
List<INT> |
|
FILTER_INT_IN_RANGE |
List<INT> |
List<INT> |
|
FILTER_LONG_GREATER_THAN |
List<LONG> |
List<LONG> |
|
FILTER_LONG_LESS_THAN |
List<LONG> |
List<LONG> |
|
FILTER_LONG_IN_RANGE |
List<LONG> |
List<LONG> |
|
FILTER_DOUBLE_GREATER_THAN |
List<DOUBLE> |
List<DOUBLE> |
|
FILTER_DOUBLE_LESS_THAN |
List<DOUBLE> |
List<DOUBLE> |
|
FILTER_DOUBLE_IN_RANGE |
List<DOUBLE> |
List<DOUBLE> |
|
FILTER_NOT_NULL |
List<T> |
List<T> |
|
FILTER_TAKE |
List<T> |
List<T> |
first n |
FILTER_SKIP |
List<T> |
List<T> |
drop first n |
FILTER_INDEX_IN_RANGE |
List<T> |
List<T> |
[from, to) |
FILTER_DISTINCT |
List<T> |
List<T> |
|
FILTER_DISTINCT_BY |
List<T> |
List<T> |
key body (T -> K); first per key, last with keepLast; a null key drops |
FILTER_TAKE_WHILE |
List<T> |
List<T> |
predicate body (T -> BOOLEAN) |
FILTER_DROP_WHILE |
List<T> |
List<T> |
predicate body (T -> BOOLEAN) |
FILTER_WHERE |
List<T> |
List<T> |
predicate body (T -> BOOLEAN) over every element; keeps each it yields true for, in order - false, null and a null element drop |
Single-element analogues of the filter family. Use them as the body of match collectors,
TakeWhile / DropWhile, FindFirst, SortBy, etc.
| Kind | Input | Output |
|---|---|---|
PREDICATE_STRING_CONTAINS |
STRING |
BOOLEAN |
PREDICATE_STRING_STARTS_WITH |
STRING |
BOOLEAN |
PREDICATE_STRING_ENDS_WITH |
STRING |
BOOLEAN |
PREDICATE_STRING_EQUALS |
STRING |
BOOLEAN |
PREDICATE_STRING_MATCHES |
STRING |
BOOLEAN |
PREDICATE_STRING_NON_EMPTY |
STRING |
BOOLEAN |
PREDICATE_INT_GREATER_THAN |
INT |
BOOLEAN |
PREDICATE_INT_LESS_THAN |
INT |
BOOLEAN |
PREDICATE_INT_IN_RANGE |
INT |
BOOLEAN |
PREDICATE_LONG_GREATER_THAN |
LONG |
BOOLEAN |
PREDICATE_LONG_LESS_THAN |
LONG |
BOOLEAN |
PREDICATE_LONG_IN_RANGE |
LONG |
BOOLEAN |
PREDICATE_DOUBLE_GREATER_THAN |
DOUBLE |
BOOLEAN |
PREDICATE_DOUBLE_LESS_THAN |
DOUBLE |
BOOLEAN |
PREDICATE_DOUBLE_IN_RANGE |
DOUBLE |
BOOLEAN |
PREDICATE_DOM_TEXT_CONTAINS |
DOM_NODE |
BOOLEAN |
PREDICATE_DOM_TEXT_MATCHES |
DOM_NODE |
BOOLEAN |
PREDICATE_DOM_HAS_ATTR |
DOM_NODE |
BOOLEAN |
PREDICATE_DOM_TAG_EQUALS |
DOM_NODE |
BOOLEAN |
PREDICATE_JSON_HAS_FIELD |
JSON_OBJECT |
BOOLEAN |
PREDICATE_JSON_FIELD_EQUALS |
JSON_OBJECT |
BOOLEAN |
PREDICATE_NOT_NULL |
T |
BOOLEAN |
PREDICATE_NOT |
BOOLEAN |
BOOLEAN |
PREDICATE_AND |
T |
BOOLEAN |
PREDICATE_OR |
T |
BOOLEAN |
PREDICATE_COMPARE |
I |
BOOLEAN |
AND / OR carry a SUB_PIPELINES_MAP of named predicate bodies and short-circuit.
COMPARE tests left OP right over two values of one input, each read by its own body
(left, right: I -> V). valueType is one of INT, LONG, FLOAT, DOUBLE, STRING or
BOOLEAN, and operator a CompareOperator named by its constant: EQUALS, NOT_EQUALS,
LESS_THAN, LESS_OR_EQUAL, GREATER_THAN, GREATER_OR_EQUAL. Numbers compare numerically
(-0.0 equals 0.0), strings by String.compareTo (case-sensitive), and booleans by equality
alone, so an ordering operator over BOOLEAN fails the load. Either body yielding null, or a
NaN on either side, yields null, and the right body does not run once the left one has yielded
null - so a FILTER_WHERE over it drops an element either side cannot be read from:
{
"kind": "FILTER_WHERE",
"elementType": "JSON_OBJECT",
"body": [{
"kind": "PREDICATE_COMPARE",
"inputType": "JSON_OBJECT",
"valueType": "INT",
"operator": "GREATER_OR_EQUAL",
"left": [{"kind": "TRANSFORM_JSON_FIELD", "fieldName": "tier"}, {"kind": "TRANSFORM_JSON_AS_INT"}],
"right": [{"kind": "TRANSFORM_CONSTANT", "inputType": "JSON_OBJECT", "outputType": "INT", "value": "3"}]
}]
}| Kind | Input | Output | Notes |
|---|---|---|---|
COLLECT_FIRST |
List<T> |
T |
|
COLLECT_LAST |
List<T> |
T |
|
COLLECT_NTH |
List<T> |
T |
element at index; past the end rejects |
COLLECT_LIST |
List<T> |
List<T> |
identity terminal marker |
COLLECT_SUB_LIST |
List<T> |
List<T> |
[from, to), to optional |
COLLECT_SET |
List<T> |
Set<T> |
|
COLLECT_JOIN |
List<STRING> |
STRING |
|
COLLECT_JSON_OBJECT_FROM_ENTRIES |
List<JSON_OBJECT> |
JSON_OBJECT |
merges {key, value} entries; the inverse of TRANSFORM_JSON_ENTRIES |
COLLECT_MAP |
I |
MAP_OUTPUT |
named bodies over one input, in declared order; a null result omits its name |
COLLECT_COUNT |
List<T> |
INT |
|
COLLECT_SUM_INT |
List<INT> |
INT |
|
COLLECT_SUM_LONG |
List<LONG> |
LONG |
|
COLLECT_SUM_DOUBLE |
List<DOUBLE> |
DOUBLE |
|
COLLECT_AVERAGE_INT |
List<INT> |
DOUBLE |
|
COLLECT_AVERAGE_LONG |
List<LONG> |
DOUBLE |
|
COLLECT_AVERAGE_DOUBLE |
List<DOUBLE> |
DOUBLE |
|
COLLECT_MIN |
List<T> |
T |
natural ordering |
COLLECT_MAX |
List<T> |
T |
natural ordering |
COLLECT_MIN_BY |
List<T> |
T |
key-extractor body (T -> K) |
COLLECT_MAX_BY |
List<T> |
T |
key-extractor body (T -> K) |
COLLECT_FIND_FIRST |
List<T> |
T |
predicate body (T -> BOOLEAN) |
COLLECT_ANY_MATCH |
List<T> |
BOOLEAN |
predicate body |
COLLECT_ALL_MATCH |
List<T> |
BOOLEAN |
predicate body |
COLLECT_NONE_MATCH |
List<T> |
BOOLEAN |
predicate body |
A stage's configuration is read off its of(...) factory: each @Configurable parameter is one
slot, keyed on the wire by the parameter's name (or @Configurable(name = ...)), and the
parameter's Java type picks the slot's FieldSpec.Type, which fixes its wire form.
FieldSpec.Type |
Factory parameter | Wire form |
|---|---|---|
STRING |
String |
string |
INT |
int / Integer |
number or numeric text; a fraction or a value past the int range fails the load |
LONG |
long / Long |
number or numeric text; a fraction or a value past the long range fails the load |
DOUBLE |
double / Double |
number or numeric text; NaN, an infinity or a number past the double range fails the load |
BOOLEAN |
boolean / Boolean |
boolean, or the text true / false in any case |
DATA_TYPE |
DataType<?> |
type label, such as "List<JSON_OBJECT>" |
SUB_PIPELINE |
List<? extends Stage<?, ?>> / Chain |
array of stages with no source - a body run against a value |
SUB_PIPELINES_MAP |
NamedChains / Map<String, List<...>> |
object of name to stage array |
TYPED_SUB_PIPELINES_MAP |
Map<String, TypedChain<?>> |
object of name to {"outputType": ..., "chain": [...]} |
PIPELINE |
DataPipeline<?> |
array of stages whose stage 0 is a source - a whole pipeline file |
STRING_MAP |
Map<String, String> |
object of strings, kept in document order |
A STRING_MAP value that is a number or boolean reads as its text; a JSON null, object or
array fails the load. TRANSFORM_VALUE_MAP's table and TRANSFORM_GROUP_BY's aggregates
are STRING_MAP slots:
{"kind": "TRANSFORM_VALUE_MAP", "table": {"PINK": "LIGHT_PURPLE", "GREY": "GRAY"}, "strict": true}A PIPELINE slot is an operand - see Reading a second document.
Every key a stage carries on the wire must be one of its slots, and every required slot must be
present - see Loading.
String json = PipelineGson.toJson(pipeline); // round-trip-safe wire format
DataPipeline<?> rebuilt = PipelineGson.fromJson(json);The wire format is a top-level JSON array of stage descriptors, each carrying a "kind"
discriminator plus its configuration fields. A stage's factory runs on load, so a body or
operand of the wrong type, or a config value the factory refuses, fails fromJson rather than
the first run.
DataPipeline outer = DataPipeline.builder()
.source(EmbedSource.of("wiki_dmg", DataTypes.INT)) // resolves at execute-time
.build();
PipelineContext ctx = PipelineContext.builder()
.withResolver(myDatabaseResolver) // host-supplied DataPipelineResolver
.build();
outer.execute(ctx);PipelineContext enforces cycle detection. Every SOURCE_EMBED enters its id before it runs, so
an embed reaching an id already running - A -> A or A -> B -> A - throws with a breadcrumb of every
active id, whether the embed is stage 0 of the pipeline or of an operand. A pipeline a host runs
directly, not through a SOURCE_EMBED, enters no id. When the host runs the very instance its
resolver answers for that id and the cycle back to it passes through an operand, the operand is
reached again before any id repeats, and the run throws the operand-cycle message instead (see
Reading a second document), which names the operand and no id.
A pipeline reads one document at stage 0. A stage that needs a second one - rows to join, a
table to look keys up in, a list to append - carries it as a pipeline operand: a PIPELINE
slot holding a whole sourced pipeline. TRANSFORM_JOIN_BY_KEY (right),
TRANSFORM_KEY_LOOKUP (table) and TRANSFORM_CONCAT (other) take one.
On the wire the operand is a stage array in the shape of a pipeline file. This joins each item row to its price row:
[
{"kind": "SOURCE_URL", "outputType": "RAW_JSON", "url": "https://example.com/items.json"},
{"kind": "PARSE_JSON"},
{"kind": "TRANSFORM_JSON_DESERIALIZE", "inputType": "JSON_ELEMENT", "outputType": "List<JSON_OBJECT>"},
{
"kind": "TRANSFORM_JOIN_BY_KEY",
"leftKey": "id",
"rightKey": "itemId",
"mode": "LEFT",
"columns": "price",
"right": [
{"kind": "SOURCE_URL", "outputType": "RAW_JSON", "url": "https://example.com/prices.json"},
{"kind": "PARSE_JSON"},
{"kind": "TRANSFORM_JSON_DESERIALIZE", "inputType": "JSON_ELEMENT", "outputType": "List<JSON_OBJECT>"}
]
}
]A saved pipeline serves as the operand through a single SOURCE_EMBED:
"right": [{"kind": "SOURCE_EMBED", "embeddedPipelineId": "prices", "outputType": "List<JSON_OBJECT>"}].
In Java the operand is any built pipeline:
DataPipeline<List<JsonObject>> prices = DataPipeline.builder()
.source(UrlSource.of(DataTypes.RAW_JSON, "https://example.com/prices.json"))
.stage(ParseJsonTransform.of())
.stage(DeserializeTransform.of(DataType.list(DataTypes.JSON_OBJECT)))
.build();
DataPipeline<List<JsonObject>> items = DataPipeline.builder()
.source(UrlSource.of(DataTypes.RAW_JSON, "https://example.com/items.json"))
.stage(ParseJsonTransform.of())
.stage(DeserializeTransform.of(DataType.list(DataTypes.JSON_OBJECT)))
.stage(JoinByKeyTransform.of("id", "itemId", "LEFT", "price", prices))
.build();- Checked when the stage is built. The owning factory validates the operand as a pipeline
file is validated, and checks that its last stage produces the type the stage consumes
(
DataPipeline.validate(DataType)). A mismatch throwsIllegalArgumentException("Invalid <Stage> operand: ..."), so a wire file with a mis-typed operand fails on load. - Evaluated at most once per run. The stage reads the operand through
PipelineContext.evaluateOperand, which runs it against the same context - fetcher, resolver, bag, tracer - the first time and holds the result, anullincluded. A lookup inside aTRANSFORM_MAPbody over a thousand elements reads its table once. The memo is keyed by the operand instance, belongs to the context, and starts empty in every newPipelineContext, so build one context per run: a context reused for a later run answers it with the operand values the first run read. The index aTRANSFORM_JOIN_BY_KEYorTRANSFORM_KEY_LOOKUPbuilds over the rows is held the same way, throughPipelineContext.derive: built once per context and released with it. - Shared, so never mutated. Every read in the run sees the same value; the stages that take an operand copy rows before changing them.
- Cycles throw. An operand whose evaluation reaches itself throws
IllegalStateException("Pipeline operand cycle detected; ..."), and an embed inside an operand goes through theSOURCE_EMBEDcycle guard. An evaluation that throws holds nothing.
Two stages go to the network, both through ctx.fetcher() - the context's UrlFetcher, with
its headers, rate limit and response cache.
SOURCE_URLfetches one configured URL at stage 0. Any status outside the2xxclass (a3xxthe fetcher does not follow, a4xxor a5xx), a transport failure, a body past the cap or a request the local rate limit refuses throws aUrlFetchExceptionand fails the run.TRANSFORM_FETCHfetches the URL its input names: the input itself, orurlTemplatewith every{}replaced by the input, which is substituted as given. The URL is sent in its ASCII form, so a character outside ASCII goes out as percent-encoded UTF-8. A space does not form a URI and a?or#ends the path, so a page name placed in a path segment has its spaces replaced first:TRANSFORM_REPLACEof" "with"_"for a MediaWiki title, or with"%20"otherwise.TRANSFORM_URL_ENCODEwrites a space as+, which a path reads as a literal plus, so it encodes a query value and not a path segment. A client error the origin answers (400to451, or a4xxthe client'sHttpStatushas no constant for outside Nginx's494-499, raised asUrlFetchException.ClientError) rejects the element withnull, so aTRANSFORM_MAPdrops a page that does not exist. A408or429is the exception: a timeout or throttling says nothing about whether the page exists, so it throws itsClientErrorand fails the run rather than silently shortening the collection. Any other status outside the2xxclass (a3xxthe fetcher does not follow, a5xx, or an Nginx444or494-499), a transport failure, a body past the cap, a local rate-limit refusal, a blank input or an input that does not form a URI throws too, so a collection is never silently short a page because the server or the network failed.
A redirect the fetcher follows is read through to the page it leads to, and a 304 answering the
fetcher's own revalidation of a cached copy is answered with the cached body. Every other 3xx
throws UrlFetchException.Redirection, which carries the code and the response headers, a
redirect's Location among them: a 300, 305 or 306, a 304 that answers no revalidation
the fetcher made (one answering an If-None-Match or If-Modified-Since among the fetcher's own
headers included), a redirect with no Location header, a redirect to another host or port from
a fetcher that sends Authorization or Cookie, or a 3xx the client's HttpStatus has no
constant for. A 2xx code HttpStatus has no constant for is read as a 200: its body is held
to the cap, passed through the guard and emitted, and the response cache does not store it. A
5xx the origin answers while the fetcher refreshes a stale cached copy is not raised when that
copy's stale-if-error window is still open as the fetch begins and no directive requires it
revalidated. The fetcher answers the cached body in its place, and the stage passes it through the
guard and emits it like a fresh one, so a finished run can hold a page an earlier fetch read. The
guard is handed the URL and the body, not the status.
Every body either stage hands on passes through the context's FetchGuard first - a body the
response cache replays included - with the URL the stage requested, the same when the fetcher
followed a redirect. The body of a status the fetch raises for never reaches it. A host installs
one once and every pipeline run against the context gets its checks, so a check specific to a
source (an error page an origin answers with a 200, say) is not repeated in each pipeline file.
A guard refuses a body by throwing, and the throw fails the run: TRANSFORM_FETCH never turns a
refusal into a dropped element, whatever the exception. A context built without one carries
FetchGuard.NOOP, which accepts every body.
PipelineContext ctx = PipelineContext.builder()
.withFetcher(fetcher)
.withFetchGuard(MyWikiGuard::check) // host-supplied: void check(URI uri, String body)
.build();Both take an optional maxBodyBytes (LONG): the largest body the fetch accepts. Absent, the
fetch is held to the fetcher's configured cap (UrlFetcherConfig, 5 MiB by default); negative,
the load fails. A body past the cap throws UrlFetchException.BodyCapExceeded, a cached body
included. An error status raises its own exception whatever the size of its body, so a 404
page larger than the cap still drops its TRANSFORM_FETCH element. One page per id, capped at
256 KiB:
[
{"kind": "SOURCE_LITERAL_LIST", "elementType": "STRING", "value": "[\"1001\", \"1002\", \"9999\"]"},
{
"kind": "TRANSFORM_MAP",
"elementInputType": "STRING",
"elementOutputType": "JSON_ELEMENT",
"body": [
{"kind": "TRANSFORM_FETCH", "outputType": "RAW_JSON", "urlTemplate": "https://api.example.com/items/{}.json", "maxBodyBytes": 262144},
{"kind": "PARSE_JSON"}
]
}
]When 9999 answers 404, the result holds the two pages that exist.
DataPipeline.Builder.build() validates the type chain eagerly and throws on any issue.
Use Builder.validate() to inspect a report on a still-under-construction builder:
ValidationReport report = DataPipeline.builder()
.source(LiteralSource.text("hi"))
.stage(ParseHtmlTransform.of()) // expects RAW_HTML, got STRING
.validate();
if (!report.isValid()) {
for (ValidationReport.Issue issue : report.issues())
System.err.println("stage #" + issue.stageIndex() + ": " + issue.message());
}Diagnostic example: Stage #1 (PARSE_HTML) expects input RAW_HTML but previous stage produced STRING.
Sub-pipeline bodies (used by Map, FlatMap, match collectors, TakeWhile / DropWhile,
Where, SortBy / MinBy / MaxBy / DistinctBy, And / Or, Compare, COLLECT_MAP,
ObjectBuild, ReplaceMatch, Coalesce, Expect, the binary arithmetic stages, Zip, Rotate
and Broadcast) are wrapped as a chain.Chain and validated the same way via
Chain.validate(seed, body, expected) when the stage is built, against the declared
element-type-in and output-type-out. A pipeline operand is validated as a whole pipeline against
the type its stage consumes (DataPipeline.validate(DataType)). Both checks run on the typed
builder path and on load alike.
A TRANSFORM_EXPECT states, in a sentence, what the value passing through it must be, and its
predicate body tests it on every run. The report lists each one before anything runs:
ValidationReport.expectations() holds every TRANSFORM_EXPECT in the pipeline, in walk order,
at any depth - a body, a named or typed body, a pipeline operand - each as an Expectation of the
top-level stage index, the path of the stage, the sentence and the type it tests. Expectations
never decide validity: a report with expectations and no issues isValid(). A SOURCE_EMBED
holds only the id of the saved pipeline it runs, which the context's resolver looks up when the
embed executes, so the report lists none of that pipeline's expectations and nothing in it marks
the embed; a host holding the resolver lists them by validating the pipeline it resolves for
EmbedSource.embeddedPipelineId().
The path follows the wire form: # and the top-level index, then per level of nesting a dot, the
slot's key, the branch name for named bodies (then .chain for a typed body) and the index in
that stage array. #1 is stage 1, #1.body[0] the first stage of its body,
#1.outputs.id.chain[1] the second stage of the id output of a
TRANSFORM_JSON_OBJECT_BUILD, and #1.right[0] stage 0 of a right operand.
for (ValidationReport.Expectation expectation : pipeline.validate().expectations())
System.out.println(expectation.path() + ": " + expectation.text());
// #1.body[0]: each id is an item idWhen a run reaches a value the body yields false or null for, it stops with
ExpectationFailedException: Expectation 'each id is an item id' failed: body yielded 'false' for input '...' of type 'STRING', the input cut to its first 120 code points. A null value
passes through without running the body, so an expectation cannot require a value to be present;
state it on the value that holds it instead - a body over each row that reads the row's id
yields null on a row with none, and that fails. A key holding JSON null is not missing:
TRANSFORM_JSON_PATH and TRANSFORM_JSON_FIELD read it as a JSON null element, and
PREDICATE_NOT_NULL and PREDICATE_JSON_HAS_FIELD both answer true for it, so such a body
passes a row whose id is JSON null. To require a non-null primitive, read it typed - every
TRANSFORM_JSON_AS_* stage yields null for JSON null, an object or an array - so a body of
TRANSFORM_JSON_PATH id, TRANSFORM_JSON_AS_STRING, PREDICATE_NOT_NULL fails on a row whose
id is missing or JSON null.
Every one of these checks asks whether the type produced is assignable to the type expected -
DataType.isAssignableTo - rather than equal to it:
JSON_OBJECTandJSON_ARRAYare assignable toJSON_ELEMENT, soTRANSFORM_JSON_STRINGIFYorTRANSFORM_JSON_PATHfollows a stage that produces an object without aTRANSFORM_JSON_DESERIALIZEin between.- A
List<X>orSet<X>is assignable to aList<Y>orSet<Y>whenXis assignable toY: rows typedList<JSON_OBJECT>feed aTRANSFORM_MAPoverJSON_ELEMENT, or aTRANSFORM_CONCAToperand onto aList<JSON_ELEMENT>. A pipeline never mutates a collection it is handed, so reading one as a collection of a wider element is safe.
Nothing is converted when the pipeline runs - a JsonObject is a JsonElement already. Nothing
else widens: RAW_HTML is not a STRING, and no numeric type is assignable to another.
PipelineGson.fromJson reads every stage strictly, at any depth - bodies and operands included -
and each refusal is an IllegalArgumentException:
- A key the stage does not declare fails the load. Every key besides
kindmust be one of the stage's slots, so a misspelt key is caught rather than read as an absent optional one:Stage 'TRANSFORM_SPLIT' does not declare key 'regx' (declared keys: [regex]). - A required key that is absent fails the load, and so does one holding JSON
null:Stage 'TRANSFORM_SPLIT' is missing required key 'regex'. - JSON
nullon an optional key reads as the key being absent, and is not written back. - An object that names a key twice fails the load, since only one of the two values would be
read:
Pipeline JSON repeats key 'regex' at '$[1].regex'. - A value of the wrong JSON shape fails the load - an object where a string belongs, say, or
anything but a stage array for a body - and so does a stage entry that is not an object or has no
kind, and aTRANSFORM_JSON_OBJECT_BUILDoutput holding a key besidesoutputTypeandchain. - A scalar that does not read as its slot's type fails the load rather than reading as
something else: a boolean slot takes a JSON boolean or the text
true/false, so"yes"or1is refused instead of read asfalse, and a number slot takes a number or numeric text:Field 'inline' holds '"yes"' but a boolean was expected. - A stage factory's refusal reaches the caller as the exception the factory threw, with its own
message:
UrlSource maxBodyBytes must not be negative but got '-1'.
v0.1, pre-release. The Stream-parity catalog (FlatMap family, terminals, match collectors, predicates, comparators, literal sources) is complete, and so is the table-shaping set built on it: joins, key lookups, concatenation and ancestor resolution across documents, group-by, distinct-by, enumerate, zip, rotate and broadcast over rows, per-element fetch, per-type arithmetic and rounding, constants and coalescing, filtering by a predicate body, comparing two values of one element, stated expectations that fail a run, and Lua, HTML-entity, JSON-escape and table-span decoding.
Open:
- No sink stage. A pipeline produces a value; writing it anywhere - a file, a table, a commit - is the host's job.
- Async / reactive
Stageexecution remains deferred. An operand or a per-element fetch runs on the calling thread.