From 8703532a8b54eaa4a75c515a88d4c5a062dea8ef Mon Sep 17 00:00:00 2001 From: John Simons Date: Fri, 21 Aug 2026 12:14:23 +1000 Subject: [PATCH] Document the ingestion pipeline and its telemetry docs/ingestion-pipeline.md covers the shape of the pipeline, how back pressure reaches the broker, what a shutdown and a hard cancellation each do to messages in flight, the three settings per instance with their ranges and what raising each one is actually for, and a table of which storage takes concurrent batches with the reason in each case. It also states the ordering rule that concurrency does not excuse, which is that the end state has to be settled by the times the events happened rather than by which batch commits last. docs/telemetry.md gains the error instance's instruments and loses the claim that PrintMetrics is how you observe error ingestion. The two instances now publish the same shape, so the sections are laid out to be read side by side and the PromQL examples say which prefix to swap. Added a short section on what the shapes mean when tuning, since the settings and the metrics that tell you how to set them now exist together. The batch duration's `full` result no longer means the transport's concurrency, it means the configured batch size, and error-ingestion-design.md no longer says a batch is what one receive cycle handed over. Three things the docs got wrong. Batch size cannot exceed the transport's concurrency, because a receive does not return until its batch has been written, so a timeout does not let a batch accumulate across receive cycles and the setting is only useful downwards. Parallel writers mean custom enrichers are called concurrently, which the settings section now says. The consecutive batch failure gauge counts whole batches failing for any reason, including announcing and forwarding, and is not the condition the custom check and the health endpoints report. --- docs/error-ingestion-design.md | 8 +- docs/ingestion-pipeline.md | 129 +++++++++++++++++++++++++++++++++ docs/telemetry.md | 48 ++++++++++-- 3 files changed, 176 insertions(+), 9 deletions(-) create mode 100644 docs/ingestion-pipeline.md diff --git a/docs/error-ingestion-design.md b/docs/error-ingestion-design.md index 38cf4a74bd..3b5e9a7963 100644 --- a/docs/error-ingestion-design.md +++ b/docs/error-ingestion-design.md @@ -8,9 +8,11 @@ for the relational persisters (PostgreSQL and SQL Server) and explains the desig are not obvious from the code, in particular why the write path is hand-written SQL rather than ordinary change-tracked entity saves. -The unit of ingestion is a **batch**: the transport hands the ingester up to `MaximumConcurrency` -messages at a time, and the whole batch is written in a single database transaction. The relevant -types are `EFIngestionUnitOfWork` (accumulation), `FailedMessageBatchWriter` (the write), and the +The unit of ingestion is a **batch**, and the whole batch is written in a single database +transaction. How batches are assembled, how many are written at once, and what can be tuned about +that is covered in [ingestion-pipeline.md](ingestion-pipeline.md); by default a batch holds up to +`MaximumConcurrency` messages and several are written concurrently. The relevant types here are +`EFIngestionUnitOfWork` (accumulation), `FailedMessageBatchWriter` (the write), and the per-provider `IFailedMessageIngestionSqlDialect` implementations (the statements that differ by provider). Retry claim insertion is a separate persistence capability behind `IRetryBatchSqlDialect`; it is used by `RetryBatchStore`, not by the ingestion unit of work. diff --git a/docs/ingestion-pipeline.md b/docs/ingestion-pipeline.md new file mode 100644 index 0000000000..2f86068966 --- /dev/null +++ b/docs/ingestion-pipeline.md @@ -0,0 +1,129 @@ +# Ingestion pipeline + +Both instances take messages off a queue, batch them, and write each batch to storage. That +machinery is one class, `IngestionPipeline` in `ServiceControl.Infrastructure`, used by +`ErrorIngestion` and `AuditIngestion`. This document covers how it behaves, what can be tuned, and +why the answer to "can batches be written in parallel" belongs to the storage rather than to the +instance. + +For what the error instance then does with a batch, see [error-ingestion-design.md](error-ingestion-design.md). +For the metrics it publishes, see [telemetry.md](telemetry.md). + +## Shape + +``` +transport receivers ──► message channel ──► batch assembler ──► batch channel ──► writers + (MaxConcurrency) (bounded) (single reader) (bounded) (1..n) +``` + +A receiver calls `Enqueue` and then waits on the message's `TaskCompletionSource`. That is what +makes the receive commit only after the message has been written: nothing is acknowledged to the +broker until the batch it landed in is in storage. A failed batch faults the completion sources of +its own messages and nothing else, so the transport redelivers exactly those. + +Both channels are bounded. When storage slows down, the batch channel fills, then the message +channel fills, and then the receivers block in `Enqueue`. Back pressure reaches the broker instead +of a queue growing in memory. + +The assembler is the only reader of the message channel, which is what keeps several writers fed +without any of them competing for messages. + +### Shutdown + +`StopAsync` stops receiving first, under the shutdown token rather than a cancelled one, so +messages already in flight finish and their receives commit. It then completes the pipeline and +waits for it to drain, and only then tears the transport infrastructure down, because batches +still forward through its dispatcher. + +A hard cancellation does not silently drop what is in flight. The assembler abandons the batch it +was building and whatever is left in the message channel, any batch no writer picked up is +abandoned, and a writer fails the batch it was holding. Every one of those receives is answered, +so they are redelivered rather than left waiting for a shutdown that has already happened. + +## Settings + +Named per instance, and read from that instance's settings root +(`ServiceControl/...` and `ServiceControl.Audit/...`). + +| Setting | Default | Range | What it does | +| --- | --- | --- | --- | +| `ErrorIngestionBatchSize` / `AuditIngestionBatchSize` | the transport's `MaximumConcurrencyLevel` | 1 to 1000 | The most messages one write handles | +| `ErrorIngestionMaxParallelWriters` / `AuditIngestionMaxParallelWriters` | the storage decides, see below | 1 to 16 | How many batches are written at once | +| `ErrorIngestionBatchTimeout` / `AuditIngestionBatchTimeout` | `00:00:00` | 0 to 5 seconds | How long a batch that is not yet full waits for more messages | + +As environment variables: + +```bash +SERVICECONTROL_ErrorIngestionBatchSize=50 +SERVICECONTROL_ErrorIngestionMaxParallelWriters=8 +SERVICECONTROL_ErrorIngestionBatchTimeout=00:00:00.100 + +SERVICECONTROL_AUDIT_AuditIngestionBatchSize=50 +SERVICECONTROL_AUDIT_AuditIngestionMaxParallelWriters=8 +SERVICECONTROL_AUDIT_AuditIngestionBatchTimeout=00:00:00.100 +``` + +### Batch size + +The default, the transport's concurrency, is also the ceiling. A receive does not return until the +batch carrying its message has been written, so the transport holds every one of its concurrency +slots open and no more than that many messages can ever be waiting. Setting a batch size above the +transport's concurrency is therefore unreachable, with or without a batch timeout: the batch never +fills, and a timeout only delays what has already arrived. Larger batches come from raising +`MaximumConcurrencyLevel`, and this setting is what holds writes below it when a storage is happier +with smaller ones. + +### Batch timeout + +Zero means a partial batch is written rather than waited on, which is what the ingestion did before +the setting existed. A non-zero value trades latency for fewer, larger writes: at volume it costs +nothing, because a full batch never waits, and at a trickle it delays each message by up to the +timeout. Start at 100ms if a storage is clearly happier with larger batches. It cannot make a batch +larger than the transport's concurrency, only fuller. + +### Parallel writers + +Only raise this for a storage whose writes are safe to interleave. Batches commit in whatever order +they finish, so this is not a free throughput knob, and the pipeline will not let you turn it on +where it is unsafe. + +More than one writer also means more than one batch is being enriched and announced at a time, so a +custom `IEnrichImportedErrorMessages` or `IEnrichImportedAuditMessages` has to be thread safe. A +single writer used to serialise them. + +## Which storages take concurrent batches + +Each ingestion unit of work factory answers `SupportsConcurrentBatches`. Where it says no, the +pipeline uses one writer whatever is configured, and logs a warning if that overrules a setting +someone actually asked for rather than a default. + +| Storage | Concurrent batches | Why | +| --- | --- | --- | +| Error, SQL Server and PostgreSQL | yes | The batch writer was built for it: upserts guarded by the attempt times so the newer attempt wins whichever transaction commits last, inserts that tolerate a competing writer's identical row, and a consistent lock order. Running several `--error-ingestion-only` hosts against one database already depends on all of it. | +| Error, RavenDB | no | Failed messages are merged by patch scripts that read and rewrite one document, and nothing orders two patches of the same document against each other. `--error-ingestion-only` refuses to start on RavenDB for the same reason. | +| Audit, RavenDB | yes | Audit documents are independent. Nothing merges two of them, and every batch gets its own bulk insert operation. | +| Audit, in memory | no | Test storage only. | + +A storage that says yes gets four writers by default. + +### The ordering that concurrency does not excuse + +Concurrent writers mean two batches touching the same message can commit in either order, which is +the same condition several ingestion hosts already create. Storage is responsible for making the +end state the same either way, and for failed messages that means comparing the times the events +happened rather than trusting arrival order: + +- a retry acknowledgement resolves a message only if no attempt newer than the retry has been + stored, because such an attempt means the message failed again afterwards +- an attempt moves the status only if it is strictly newer than the newest attempt already stored, + so a redelivery of the attempt already there cannot undo a resolve or an archive + +## Where things live + +| | | +| --- | --- | +| `ServiceControl.Infrastructure/Ingestion/IngestionPipeline.cs` | The channels, the assembler and the writers | +| `ServiceControl.Infrastructure/Ingestion/IngestionSettingsReader.cs` | Reading, validating and resolving the settings above | +| `ServiceControl.Infrastructure/Ingestion/Metrics/` | The metric scopes both instances report through | +| `ServiceControl/Operations/ErrorIngestion.cs` | Transport, watchdog and fault policy for the error queue | +| `ServiceControl.Audit/Auditing/AuditIngestion.cs` | The same for the audit queue | diff --git a/docs/telemetry.md b/docs/telemetry.md index 2daded1bf4..5dbad512d9 100644 --- a/docs/telemetry.md +++ b/docs/telemetry.md @@ -2,18 +2,37 @@ Instances can be configured to emit telemetry to aid in performance testing or troubleshooting performance-related issues. +Both the error and the audit instance report their ingestion the same way. Set `OtlpEndpointUrl` +in that instance's settings root to a valid [OTLP endpoint url](https://opentelemetry.io/docs/specs/otel/protocol/exporter/#configuration-options). +Only GRPC endpoints are supported at this stage. + +The instruments differ only in their prefix and in the categories a message can fall into, so the +same dashboard works for both with the prefix swapped. What the batches being measured actually are +is covered in [ingestion-pipeline.md](ingestion-pipeline.md). + ## Error -Setting `ServiceControl/PrintMetrics` to `true` will print metrics to the logs at `INFO` level. +Meter `Particular.ServiceControl`, configured with `ServiceControl/OtlpEndpointUrl`. -## Audit +- `sc.error.ingestion.batch_duration_seconds` - Message batch processing duration in seconds + - `result` - Whether the batch was written at its configured size: `full`, `partial` or `failed` +- `sc.error.ingestion.storage_duration_seconds` - The storage write of a batch, without the announcing and forwarding around it +- `sc.error.ingestion.message_duration_seconds` - Error message processing duration in seconds + - `message.category` - What the message is: `failed-message` or `retry-confirmation` + - `result` - The outcome: `success`, `failed` or `skipped` (if the message was filtered out and skipped) +- `sc.error.ingestion.failures_total` - Failure counter + - `message.category` - What the message is: `failed-message` or `retry-confirmation` + - `result` - How the failure was resolved: `retry` or `stored-poison` +- `sc.error.ingestion.consecutive_batch_failures_total` - Consecutive batch failures -Set `ServiceControl.Audit/OtlpEndpointUrl` to a valid [OTLP endpoint url](https://opentelemetry.io/docs/specs/otel/protocol/exporter/#configuration-options). Only GRPC endpoints are supported at this stage. +`ServiceControl/PrintMetrics` predates this and no longer has anything to print. -The following ingestion metrics with their corresponding dimensions are available: +## Audit + +Meter `Particular.ServiceControl.Audit`, configured with `ServiceControl.Audit/OtlpEndpointUrl`. - `sc.audit.ingestion.batch_duration_seconds` - Message batch processing duration in seconds - - `result` - Indicates if the full batch size was used (batch size == max concurrency of the transport): `full`, `partial` or `failed` + - `result` - Whether the batch was written at its configured size: `full`, `partial` or `failed` - `sc.audit.ingestion.message_duration_seconds` - Audit message processing duration in seconds - `message.category` - Indicates the category of the message ingested: `audit-message`, `saga-update` or `control-message` - `result` - Indicates the outcome of the operation: `success`, `failed` or `skipped` (if the message was filtered out and skipped) @@ -22,12 +41,29 @@ The following ingestion metrics with their corresponding dimensions are availabl - `result` - Indicates how the failure was resolved: `retry` or `stored-poison` - `sc.audit.ingestion.consecutive_batch_failures_total` - Consecutive batch failures -Example queries in PromQL: +## Reading the ingestion metrics + +Example queries in PromQL, shown for audit. Swap `sc_audit` for `sc_error` for the error instance. - Ingestion rate: `sum (rate(sc_audit_ingestion_message_duration_seconds_count[5m])) by (exported_job)` - Failure rate: `sum(rate(sc_audit_ingestion_failures_total[5m])) by (exported_job,result)` - Message duration: `histogram_quantile(0.9,sum(rate(sc_audit_ingestion_message_duration_seconds_bucket[5m])) by (le,exported_job))` +What the shapes mean when tuning: + +- `result="full"` dominating the batch duration histogram means batches are being filled, so the + limit is the write and not the arrival of messages. More writers, or a larger batch size, is + where to look. +- `result="partial"` dominating at high throughput means the pipeline is writing before batches + fill. A batch timeout is what makes them accumulate. +- Message duration far above batch duration means messages are queueing behind the writers rather + than being slow to write. +- `consecutive_batch_failures_total` counts whole batches that failed in a row, whatever ended them: + the storage write, but equally announcing or forwarding. It is not a tuning signal, and it is not + the same thing the ingestion custom check and the health endpoints report, which is the critical + failure the watchdog acts on. A gauge climbing while health stays clear means batches are failing + and being retried. + Example Grafana dashboard - https://github.com/andreasohlund/Docker/blob/main/otel-monitoring/grafana-platform-template.json ## Monitoring