Skip to main content

Sample Node

The sample node keeps one randomly picked record out of every N records matched by a rule and drops the rest. Rules are checked in order and the first match wins; each rule has its own sample rate. Every kept record carries a field (by default __sampled__) set to the number of records it stands for. Records that no rule matches pass through unchanged.

For more details of using FEEL, please refer to FEEL Reference

Key Features​

  • Ordered rules: Each record goes to the first rule whose condition is true; a rule without a condition matches every record
  • Uniform 1-in-N sampling: Each rule keeps one record picked uniformly at random from every group of sampleRate records it matches
  • Rate tagging: Kept records carry the size of the group they stand for, so once input has ended, the tags add up to the number of matched records
  • End-of-input flush: When input ends, each rule's incomplete group still keeps one record, tagged with the group's actual size
  • Pass-through: Records that match no rule are passed through unchanged

Config Object​

FieldTypeRequiredDefaultDescription
rulesList<Rule>Yes—Ordered list of rules. Must not be empty.
sampleRateFieldStringNo__sampled__Top-level field set on every kept record. Must not be blank. Used as a literal key, not a path.

Each Rule:

FieldTypeRequiredDefaultDescription
conditionStringNo—FEEL expression. The rule matches a record when the expression returns true. Omit to match every record. Must not be blank.
sampleRateIntegerYes—Keep one record out of every sampleRate matched records. Between 1 and 2147483647.

sampleRate must be a whole number: 3.0, "3" and out-of-range values are rejected. Unknown fields are rejected at both levels, so a misspelled field name fails validation instead of being ignored. A condition that does not compile fails validation with the rule's index, for example invalid 'rules[0].condition'.

How Sampling Works​

Rule Matching​

Rules are checked in order and a record goes to the first rule that matches. A rule matches when its condition returns the boolean true, or always when it has no condition. A condition that returns any other value, or throws, does not match, and the next rule is checked. Records that match no rule are passed through unchanged, without the rate field. A rule without a condition matches every record, so rules after it never match.

Groups of N​

Each rule counts the records it matches in groups of sampleRate. Within a group, the node keeps one candidate, replacing it so that every record of the group is equally likely to be the one kept (reservoir sampling). When the sampleRate-th record of the group arrives, the candidate is emitted and a new group starts. Every other record of the group is dropped.

The kept record is emitted when its group completes, not when it arrives, so output order can differ from input order. With sampleRate: 1 every matched record is emitted, tagged with 1.

The Rate Field​

Every kept record is emitted with a top-level field named by sampleRateField set to the size of its group: the sampleRate, or the actual size for a group flushed at end of input. If the record already has that field, it is overwritten. The name is not a path: sample.rate creates a field named sample.rate.

Over a whole input, the rate field of the kept records adds up to the number of records the rules matched.

End of Input​

When the input ends, each rule with an incomplete group emits its candidate, tagged with the number of records in that group. Input ends when a finite source reaches its end or when the source stops because of an error. Test runs in the Fleak Workflow Builder also end the input.

Stopping the process, for example with Ctrl-C or SIGTERM, does not end the input this way: each rule's incomplete group is lost, with up to sampleRate - 1 matched records that are never emitted or counted in the rate field.

Scope of State​

Groups live in the memory of each running pipeline instance. They are not shared across instances, and each instance samples the records it receives.

Metrics​

In addition to the standard input_event_count, output_event_count and error_event_count counters, the node emits sample_dropped_count for each dropped record.

Examples​

Per-Condition Rates​

Keep 1 in 100 health checks and 1 in 10 successful requests, and pass all remaining records through unchanged:

{
"rules": [
{ "condition": "$.path == \"/healthz\"", "sampleRate": 100 },
{ "condition": "$.status == 200", "sampleRate": 10 }
],
"sampleRateField": "sample_rate"
}

Records that match neither rule pass through without sample_rate.

Sample Input and Output​

With these rules:

{
"rules": [
{ "condition": "$.status == 200", "sampleRate": 3 },
{ "sampleRate": 2 }
]
}
OrderInputRuleOutput
1{"id": "a1", "status": 200}rules[0]—
2{"id": "a2", "status": 200}rules[0]—
3{"id": "a3", "status": 200}rules[0]one of a1–a3, e.g. {"id": "a2", "status": 200, "__sampled__": 3}
4{"id": "a4", "status": 200}rules[0]—
5{"id": "b1", "status": 500}rules[1]—
6{"id": "b2", "status": 500}rules[1]one of b1–b2, e.g. {"id": "b1", "status": 500, "__sampled__": 2}
7{"id": "b3", "status": 500}rules[1]—
end{"id": "a4", "status": 200, "__sampled__": 1}, {"id": "b3", "status": 500, "__sampled__": 1}

DAG Definition​

jobContext:
metricTags: {}
otherProperties: {}

dag:
- id: "source"
commandName: "kafkasource"
config:
broker: "localhost:9092"
topic: "access-logs"
groupId: "access-log-sampler"
encodingType: "JSON_OBJECT"
outputs:
- "sample"

- id: "sample"
commandName: "sample"
config:
rules:
- condition: "$.status == 200"
sampleRate: 10
outputs:
- "sink"

- id: "sink"
commandName: "stdout"
config:
encodingType: "JSON_OBJECT"
outputs: []
  • throttle: Limit records per key in each time period instead of sampling
  • filter: Drop records that don't match a condition
  • eval: Compute a field to sample on