k8s-infra-cluster/infra/benthos/templates/poll_template_id_serial.yaml
2026-03-25 11:46:39 +01:00

60 lines
1.6 KiB
YAML

# templates/poll_template.yaml
name: sql_poller
type: input
# using BIGSERIAL
fields:
- name: table_name
type: string
- name: dsn
type: string
- name: id_column
type: string
default: "id"
mapping: |
# 1. The base input is a generator that ticks every 2s
root.generate.interval = "2s"
root.generate.mapping = "root = {}"
# 2. We attach the processing logic directly to this input
root.processors = [
{
"cache": {
"resource": "yugabyte_checkpoints",
"operator": "get",
"key": "last_id_" + this.table_name
}
},
{ "catch": [ { "mapping": "root = \"0\"" } ] }, # Default to 0 if cache missing
{ "mapping": "meta last_id = content().string()" },
# 3. Perform the Query
{
"sql_raw": {
"driver": "postgres",
"dsn": this.dsn,
"query": "SELECT * FROM " + this.table_name + " WHERE " + this.id_column + " > $1 ORDER BY " + this.id_column + " ASC LIMIT 500;",
"args_mapping": "root = [ meta(\"last_id\").number() ]"
}
},
# 4. Filter empty results
{
"mapping": "root = if this.type() != \"array\" || this.length() == 0 { deleted() } else { this }"
},
# 5. Update the Cache with the new max ID
{
"branch": {
"processors": [
{ "mapping": "root = this.index(-1)." + this.id_column + ".string()" },
{ "cache": { "resource": "yugabyte_checkpoints", "operator": "set", "key": "last_id_" + this.table_name, "value": "${! content() }" }}
]
}
},
# 6. Flatten the array into individual messages
{ "unarchive": { "format": "json_array" } }
]