Down Sampling
Down sampling aggregates incoming records at the source layer before payload decoding and transformation.
You can also aggregate streaming data using time windows. However, time windows decode every incoming event before evaluating the window condition. Discarded events consume CPU cycles during decoding.
Source down sampling filters or merges raw events before decoding. This design reduces CPU usage and memory consumption. As the ratio between input frequency and sampling frequency increases, efficiency gains become more significant.
Applicable Source Categories
Connectors belong to two categories based on data ingestion:
- Pull Sources: The engine polls the external system at scheduled intervals. Set the polling interval to control down sampling.
- Push Sources: The external publisher controls ingestion frequency. The push source buffers incoming records and emits them based on the configured down sampling strategy.
Configuration Properties
Configure down sampling using these source properties:
interval
Specifies the sampling duration using a duration string (such as "10s" or "500ms").
- For pull sources,
intervaldefines the polling period. - For push sources, the connector accumulates incoming events during the interval and emits records when the interval elapses.
mergeField
Defines the column name used for record aggregation (for example, "id").
The default down sampling strategy emits the last record received during the sampling period. If you configure mergeField, the engine retains the last record for each distinct key value and merges them into a single composite record.
NOTE
mergeFieldrequires formats that support partial decoding (such as JSON). Custom formats can implementmessage.PartialDecoder.mergeFieldsupports top-level non-composite fields. Nested objects or array structures resolve to the same key.
Down Sampling Strategies
Source down sampling converts multiple events received during an interval into a single output event. rekuiper supports two strategies:
- Latest Value: Emits the final record received during the sampling period.
- Column Aggregation: Merges the latest records for each unique key in
mergeField.
Strategy 1: Latest Value
This strategy outputs the last record received during the sampling interval.
Step 1: Create the Source Configuration
Create an MQTT configuration named onesec with a 1-second sampling period:
PUT http://{{host}}/metadata/sources/mqtt/confKeys/onesec
Content-Type: application/json
{
"interval": "1s"
}Step 2: Create the Stream
Create a stream that uses the onesec configuration:
POST http://{{host}}/streams
Content-Type: application/json
{
"sql": "CREATE STREAM mqttOneSec() WITH (TYPE=\"mqtt\", FORMAT=\"json\", DATASOURCE=\"demo\", CONF_KEY=\"onesec\");"
}Step 3: Create the Rule
Create a rule to process the downsampled stream:
POST http://{{host}}/rules
Content-Type: application/json
{
"id": "ruleOneSecLatest",
"sql": "SELECT * FROM mqttOneSec",
"actions": [
{
"mqtt": {
"server": "tcp://127.0.0.1:1883",
"topic": "result/onesec",
"sendSingle": true
}
}
]
}The rule receives records once per second and emits the latest record for each interval.
Strategy 2: Column Aggregation
This strategy groups records by key and combines the latest values across distinct fields into a single record.
Consider this sequence of input events within a 1-second window:
{"id": 1, "temperature": 20}
{"id": 2, "humidity": 80}
{"id": 1, "temperature": 30}The engine merges the events into one composite record:
{
"id": 1,
"temperature": 30,
"humidity": 80
}Step 1: Create the Merge Configuration
Create an MQTT configuration named onesec_merge with interval and mergeField:
PUT http://{{host}}/metadata/sources/mqtt/confKeys/onesec_merge
Content-Type: application/json
{
"interval": "1s",
"mergeField": "id"
}Step 2: Create the Stream
Create a stream referencing onesec_merge:
POST http://{{host}}/streams
Content-Type: application/json
{
"sql": "CREATE STREAM mqttOneSecM() WITH (TYPE=\"mqtt\", FORMAT=\"json\", DATASOURCE=\"demo\", CONF_KEY=\"onesec_merge\");"
}Step 3: Create the Rule
Create a rule to process the merged stream:
POST http://{{host}}/rules
Content-Type: application/json
{
"id": "RuleOneSecM",
"sql": "SELECT * FROM mqttOneSecM",
"actions": [
{
"mqtt": {
"server": "tcp://127.0.0.1:1883",
"topic": "result/onesecm",
"sendSingle": true
}
}
]
}Full Aggregation with Windows
To aggregate records across all columns without specifying a key field, use a time window with the merge_agg function:
SELECT merge_agg(*)
FROM normalStream
GROUP BY TumblingWindow(ss, 1);Source down sampling minimizes decoding overhead by inspecting only key fields before decoding. Full aggregation decodes every event payload. For full aggregation, ingest records using a standard stream and apply window functions in SQL.
Down Sampling Observability
The ratelimit operator executes source down sampling. Inspect the ratelimit metrics to monitor down sampling efficiency.
In the following status output:
source_mqttOneMiMerge_0_records_out_totalindicates the MQTT connector ingested 25 records.op_2_ratelimit_0_records_in_totalis 25, andop_2_ratelimit_0_records_out_totalis 1. The operator reduced 25 incoming records to 1 record.- The downstream
op_3_payload_decoder_0decoded only the single downsampled record.
{
"status": "running",
"lastStartTimestamp": "1720151899579",
"lastStopTimestamp": "0",
"nextStopTimestamp": "0",
"source_mqttOneMiMerge_0_records_in_total": 25,
"source_mqttOneMiMerge_0_records_out_total": 25,
"source_mqttOneMiMerge_0_messages_processed_total": 25,
"source_mqttOneMiMerge_0_process_latency_us": 0,
"source_mqttOneMiMerge_0_buffer_length": 0,
"source_mqttOneMiMerge_0_last_invocation": "2024-07-05T11:58:40.733398",
"source_mqttOneMiMerge_0_exceptions_total": 0,
"source_mqttOneMiMerge_0_last_exception": "",
"source_mqttOneMiMerge_0_last_exception_time": 0,
"op_2_ratelimit_0_records_in_total": 25,
"op_2_ratelimit_0_records_out_total": 1,
"op_2_ratelimit_0_messages_processed_total": 25,
"op_2_ratelimit_0_process_latency_us": 0,
"op_2_ratelimit_0_buffer_length": 0,
"op_2_ratelimit_0_last_invocation": "2024-07-05T11:58:40.733398",
"op_2_ratelimit_0_exceptions_total": 0,
"op_2_ratelimit_0_last_exception": "",
"op_2_ratelimit_0_last_exception_time": 0,
"op_3_payload_decoder_0_records_in_total": 1,
"op_3_payload_decoder_0_records_out_total": 1,
"op_3_payload_decoder_0_messages_processed_total": 1,
"op_3_payload_decoder_0_process_latency_us": 0,
"op_3_payload_decoder_0_buffer_length": 0,
"op_3_payload_decoder_0_last_invocation": "2024-07-05T11:59:19.59698",
"op_3_payload_decoder_0_exceptions_total": 0,
"op_3_payload_decoder_0_last_exception": "",
"op_3_payload_decoder_0_last_exception_time": 0,
"op_4_project_0_records_in_total": 1,
"op_4_project_0_records_out_total": 1,
"op_4_project_0_messages_processed_total": 1,
"op_4_project_0_process_latency_us": 0,
"op_4_project_0_buffer_length": 0,
"op_4_project_0_last_invocation": "2024-07-05T11:59:19.59698",
"op_4_project_0_exceptions_total": 0,
"op_4_project_0_last_exception": "",
"op_4_project_0_last_exception_time": 0,
"op_mqtt_0_0_transform_0_records_in_total": 1,
"op_mqtt_0_0_transform_0_records_out_total": 1,
"op_mqtt_0_0_transform_0_messages_processed_total": 1,
"op_mqtt_0_0_transform_0_process_latency_us": 0,
"op_mqtt_0_0_transform_0_buffer_length": 0,
"op_mqtt_0_0_transform_0_last_invocation": "2024-07-05T11:59:19.59698",
"op_mqtt_0_0_transform_0_exceptions_total": 0,
"op_mqtt_0_0_transform_0_last_exception": "",
"op_mqtt_0_0_transform_0_last_exception_time": 0,
"op_mqtt_0_1_encode_0_records_in_total": 1,
"op_mqtt_0_1_encode_0_records_out_total": 1,
"op_mqtt_0_1_encode_0_messages_processed_total": 1,
"op_mqtt_0_1_encode_0_process_latency_us": 0,
"op_mqtt_0_1_encode_0_buffer_length": 0,
"op_mqtt_0_1_encode_0_last_invocation": "2024-07-05T11:59:19.59698",
"op_mqtt_0_1_encode_0_exceptions_total": 0,
"op_mqtt_0_1_encode_0_last_exception": "",
"op_mqtt_0_1_encode_0_last_exception_time": 0,
"sink_mqtt_0_0_records_in_total": 1,
"sink_mqtt_0_0_records_out_total": 1,
"sink_mqtt_0_0_messages_processed_total": 1,
"sink_mqtt_0_0_process_latency_us": 0,
"sink_mqtt_0_0_buffer_length": 0,
"sink_mqtt_0_0_last_invocation": "2024-07-05T11:59:19.59698",
"sink_mqtt_0_0_exceptions_total": 0,
"sink_mqtt_0_0_last_exception": "",
"sink_mqtt_0_0_last_exception_time": 0
}