Beats and Elastic Agent JSON in, DFE-normalised events out.
dfe-transform-elastic is a batch transform service. It consumes JSON event batches produced by filebeat, winlogbeat, metricbeat, auditbeat, heartbeat, packetbeat or Elastic Agent, from Kafka or over a direct gRPC push, applies the Elastic ingest-pipeline logic for that source, and produces normalised events downstream the same way.
The transform logic is compiled in. There is no interpreter, no scripting VM and no plugin system: each supported source is a Rust module, and the service resolves one of them by name at startup. That is the whole design decision, and everything else follows from it -- the container is just the binary, the hot path has no dynamic dispatch per event, and adding a source means a release rather than a config change.
It is built on the scalo data-plane runtime: config cascade, logging, metrics, the Kafka and gRPC transports, health probes, memory guard and scaling pressure.
Sibling services. dfe-transform-vrl runs user-supplied VRL programs and
is the right choice when the transform must be changed without a release. This
service is the right choice when the transform is a known Elastic pipeline and
throughput matters.
cargo build --release
./target/release/dfe-transform-elastic --config config.yamlList the sources a build can transform:
./target/release/dfe-transform-elastic sourcesLoaded from an explicit --config path, or from scalo's config cascade under
the DFE_TRANSFORM_ELASTIC environment prefix.
source:
name: filebeat.okta.default # one of `sources` above
brokers: ["kafka:9092"]
topics: ["raw_events"]
group_id: dfe-transform-elastic-my-pipeline
batch_size: 20000
sink:
topic: normalised_events
brokers: ["kafka:9092"] # defaults to the source brokers
max_message_bytes: 15728640 # ceiling on one outbound record (15 MiB)That is the shape, not the whole surface. config.example.yaml is the COMPLETE
set of defaults with every key commented, generated from the deployment contract
and pinned against drift by a test -- read it rather than this snippet when you
need a key that is not here. dfe-transform-elastic emit-config reprints it.
Every transformed event goes out as its OWN record, because dfe-loader parses
one JSON document per message. max_message_bytes bounds each one, so keep it
below your broker's message.max.bytes.
A config naming a source this build does not carry is rejected at startup, not discovered at the first batch.
Events arrive as NDJSON, one JSON object per line, wrapped by one of three
producers: Beats and Elastic Agent (beats), dfe-receiver (receiver) or
dfe-fetcher (fetcher). source.envelope defaults to auto, which reads the
family off each event, and naming one pins it. The producer's own field names
are stripped once they have been lifted onto ECS, except the one dfe-loader
routes on -- _source from the receiver and _source_fetcher from the fetcher
come through unchanged. What each family carries, and
what a payload with no marker does, are in
docs/architecture.md.
JSON is the only payload format. MessagePack, supported in DFE/XDR 2.0 and 2.1, is deprecated in DFE 2.2 and no longer accepted: the JSON path (SIMD parsing with sonic-rs, zstd on the wire) is fast enough that MessagePack gave no CPU saving.
Each side also carries a transport: bus, the default, is Kafka, and
direct is a scalo Push listener inbound and a gRPC push outbound. Where the
values come from, which env spelling reaches a --config file, and which
sections that file cannot carry are in
docs/configuration.md.
The service reads whatever is on the topic, so its failure modes are stated rather than assumed:
- Invalid UTF-8 is decoded with U+FFFD replacements, matching what Beats
itself substitutes for a file it cannot decode. The payload survives; the
substitution is counted on
lossy_payloads_total. - Input that is not JSON is refused and dead-lettered with the reason
payload is not JSON, and counted onparse_errors_total. A record no line of which parses -- MessagePack, or any other binary -- is refused whole, with the bytes the producer sent. In a record that is otherwise JSON, only the bad line is refused and the events beside it still go. This service configures no DLQ, so the dead letter is dropped and counted onpipeline_dead_letters_dropped_total{reason="dead_letter"}, and its record's source is released. - An event whose transform errors is counted on
events_errored_totaland left out of the output. The batch continues. - An event too large for one Kafka record is dropped and counted on
events_oversize_total. No broker would accept it, and retrying it forever would block the partition behind it. - A sink that refuses for now -- a full producer queue, a broker outage -- is retried for as long as it refuses, counted on
send_backpressure_total. The Kafka offset commit, or the answer to a push, is held until the sink takes the block, so an outage is waited out and nothing after the block is committed past it. - A send no retry can fix STOPS the service with the block unreleased, and the restarted consumer reads it again. Delivery is at-least-once, so a downstream consumer must be idempotent.
- A record the sink transport would refuse is dropped before the send and counted on
pipeline_dead_letters_dropped_total, never counted delivered.
Non-English text is a tested case, not an edge case. tests/unicode.rs
runs every registered transform against seventeen scripts and a set of
degenerate inputs.
The container image, Helm chart, compose fragment and KEDA scaler are all
generated from one deployment contract in src/deployment.rs:
dfe-transform-elastic emit-dockerfile > Dockerfile
dfe-transform-elastic emit-chart chart/dfe-transform-elastic
dfe-transform-elastic emit-compose
dfe-transform-elastic generate-artefacts --output-dir docsgenerate-artefacts writes into docs/: metrics-manifest.json,
deployment-contract.json, container-manifest.json, Dockerfile.runtime,
and the reflectable config pair config-schema.* and capability-catalog.*.
It writes an argocd-application.yaml only when deployment.argocd.repo_url is set in the config cascade, and this repo sets none. scalo writes that Application's source as path: chart, the chart here is chart/dfe-transform-elastic, and scalo has no setting for the path, so the Application would sync nothing. The file is gitignored so a local cascade that sets the URL cannot commit one.
Not every artefact is pinned against a fresh regen. committed_config_artefacts_do_not_drift in src/deployment.rs compares config-schema.{json,yaml} and capability-catalog.{json,yaml}, and its sibling tests do the same for the committed deployment-contract.json, Dockerfile, config.example.yaml and the chart's config: block. Nothing compares container-manifest.json or Dockerfile.runtime, so those can and do fall behind -- re-run the command when src/deployment.rs changes rather than assuming a test caught it.
metrics-manifest.json is compared by metric NAME SET in src/metrics.rs, not byte for byte, because it also records the version and commit of the build that wrote it.
The image expects the release binary in the build context:
cargo build --release
cp target/release/dfe-transform-elastic .
docker build -t dfe-transform-elastic .The source catalogue, sources.yaml, is compiled into the binary. dfe-transform-elastic emit-catalogue prints it, so a deployment takes the catalogue dfe-engine offers from the image it already pulls, with no credential for this repo.
/livez,/readyz,/metricsand/metrics/manifest, all served from the one listener on port 9090. There is no second health port./scaling/pressureserves a single weighted figure to KEDA: consumer lag at 0.70, batch saturation at 0.30, with memory as a hard gate that forces pressure to 100 before an OOM.dfe-transform-elastic metrics-manifestprints the full metric catalogue without starting the service. The committed copy isdocs/metrics-manifest.json.
cargo test --workspace --all-features
hyperi-ci checklibrdkafka 2.12.1 or later is needed for the kafka feature, which is on by
default. Building without it still works:
cargo build --no-default-featuresdocs/ carries everything deeper than this page -- architecture.md for the code map, parity.md for what the service promises against Elastic's own output, and compat.md for working a source towards it.
BUSL-1.1. See LICENSE, and COMMERCIAL.md for commercial terms.