|
| 1 | +# Ingestion pipeline |
| 2 | + |
| 3 | +Both instances take messages off a queue, batch them, and write each batch to storage. That |
| 4 | +machinery is one class, `IngestionPipeline` in `ServiceControl.Infrastructure`, used by |
| 5 | +`ErrorIngestion` and `AuditIngestion`. This document covers how it behaves, what can be tuned, and |
| 6 | +why the answer to "can batches be written in parallel" belongs to the storage rather than to the |
| 7 | +instance. |
| 8 | + |
| 9 | +For what the error instance then does with a batch, see [error-ingestion-design.md](error-ingestion-design.md). |
| 10 | +For the metrics it publishes, see [telemetry.md](telemetry.md). |
| 11 | + |
| 12 | +## Shape |
| 13 | + |
| 14 | +``` |
| 15 | +transport receivers ──► message channel ──► batch assembler ──► batch channel ──► writers |
| 16 | + (MaxConcurrency) (bounded) (single reader) (bounded) (1..n) |
| 17 | +``` |
| 18 | + |
| 19 | +A receiver calls `Enqueue` and then waits on the message's `TaskCompletionSource`. That is what |
| 20 | +makes the receive commit only after the message has been written: nothing is acknowledged to the |
| 21 | +broker until the batch it landed in is in storage. A failed batch faults the completion sources of |
| 22 | +its own messages and nothing else, so the transport redelivers exactly those. |
| 23 | + |
| 24 | +Both channels are bounded. When storage slows down, the batch channel fills, then the message |
| 25 | +channel fills, and then the receivers block in `Enqueue`. Back pressure reaches the broker instead |
| 26 | +of a queue growing in memory. |
| 27 | + |
| 28 | +The assembler is the only reader of the message channel, which is what keeps several writers fed |
| 29 | +without any of them competing for messages. |
| 30 | + |
| 31 | +### Shutdown |
| 32 | + |
| 33 | +`StopAsync` stops receiving first, under the shutdown token rather than a cancelled one, so |
| 34 | +messages already in flight finish and their receives commit. It then completes the pipeline and |
| 35 | +waits for it to drain, and only then tears the transport infrastructure down, because batches |
| 36 | +still forward through its dispatcher. |
| 37 | + |
| 38 | +A hard cancellation does not silently drop what is in flight. The assembler abandons the batch it |
| 39 | +was building and whatever is left in the message channel, any batch no writer picked up is |
| 40 | +abandoned, and a writer fails the batch it was holding. Every one of those receives is answered, |
| 41 | +so they are redelivered rather than left waiting for a shutdown that has already happened. |
| 42 | + |
| 43 | +## Settings |
| 44 | + |
| 45 | +Named per instance, and read from that instance's settings root |
| 46 | +(`ServiceControl/...` and `ServiceControl.Audit/...`). |
| 47 | + |
| 48 | +| Setting | Default | Range | What it does | |
| 49 | +| --- | --- | --- | --- | |
| 50 | +| `ErrorIngestionBatchSize` / `AuditIngestionBatchSize` | the transport's `MaximumConcurrencyLevel` | 1 to 1000 | The most messages one write handles | |
| 51 | +| `ErrorIngestionMaxParallelWriters` / `AuditIngestionMaxParallelWriters` | the storage decides, see below | 1 to 16 | How many batches are written at once | |
| 52 | +| `ErrorIngestionBatchTimeout` / `AuditIngestionBatchTimeout` | `00:00:00` | 0 to 5 seconds | How long a batch that is not yet full waits for more messages | |
| 53 | + |
| 54 | +As environment variables: |
| 55 | + |
| 56 | +```bash |
| 57 | +SERVICECONTROL_ErrorIngestionBatchSize=50 |
| 58 | +SERVICECONTROL_ErrorIngestionMaxParallelWriters=8 |
| 59 | +SERVICECONTROL_ErrorIngestionBatchTimeout=00:00:00.100 |
| 60 | + |
| 61 | +SERVICECONTROL_AUDIT_AuditIngestionBatchSize=50 |
| 62 | +SERVICECONTROL_AUDIT_AuditIngestionMaxParallelWriters=8 |
| 63 | +SERVICECONTROL_AUDIT_AuditIngestionBatchTimeout=00:00:00.100 |
| 64 | +``` |
| 65 | + |
| 66 | +### Batch size |
| 67 | + |
| 68 | +The default, the transport's concurrency, is also the ceiling. A receive does not return until the |
| 69 | +batch carrying its message has been written, so the transport holds every one of its concurrency |
| 70 | +slots open and no more than that many messages can ever be waiting. Setting a batch size above the |
| 71 | +transport's concurrency is therefore unreachable, with or without a batch timeout: the batch never |
| 72 | +fills, and a timeout only delays what has already arrived. Larger batches come from raising |
| 73 | +`MaximumConcurrencyLevel`, and this setting is what holds writes below it when a storage is happier |
| 74 | +with smaller ones. |
| 75 | + |
| 76 | +### Batch timeout |
| 77 | + |
| 78 | +Zero means a partial batch is written rather than waited on, which is what the ingestion did before |
| 79 | +the setting existed. A non-zero value trades latency for fewer, larger writes: at volume it costs |
| 80 | +nothing, because a full batch never waits, and at a trickle it delays each message by up to the |
| 81 | +timeout. Start at 100ms if a storage is clearly happier with larger batches. It cannot make a batch |
| 82 | +larger than the transport's concurrency, only fuller. |
| 83 | + |
| 84 | +### Parallel writers |
| 85 | + |
| 86 | +Only raise this for a storage whose writes are safe to interleave. Batches commit in whatever order |
| 87 | +they finish, so this is not a free throughput knob, and the pipeline will not let you turn it on |
| 88 | +where it is unsafe. |
| 89 | + |
| 90 | +More than one writer also means more than one batch is being enriched and announced at a time, so a |
| 91 | +custom `IEnrichImportedErrorMessages` or `IEnrichImportedAuditMessages` has to be thread safe. A |
| 92 | +single writer used to serialise them. |
| 93 | + |
| 94 | +## Which storages take concurrent batches |
| 95 | + |
| 96 | +Each ingestion unit of work factory answers `SupportsConcurrentBatches`. Where it says no, the |
| 97 | +pipeline uses one writer whatever is configured, and logs a warning if that overrules a setting |
| 98 | +someone actually asked for rather than a default. |
| 99 | + |
| 100 | +| Storage | Concurrent batches | Why | |
| 101 | +| --- | --- | --- | |
| 102 | +| 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. | |
| 103 | +| 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. | |
| 104 | +| Audit, RavenDB | yes | Audit documents are independent. Nothing merges two of them, and every batch gets its own bulk insert operation. | |
| 105 | +| Audit, in memory | no | Test storage only. | |
| 106 | + |
| 107 | +A storage that says yes gets four writers by default. |
| 108 | + |
| 109 | +### The ordering that concurrency does not excuse |
| 110 | + |
| 111 | +Concurrent writers mean two batches touching the same message can commit in either order, which is |
| 112 | +the same condition several ingestion hosts already create. Storage is responsible for making the |
| 113 | +end state the same either way, and for failed messages that means comparing the times the events |
| 114 | +happened rather than trusting arrival order: |
| 115 | + |
| 116 | +- a retry acknowledgement resolves a message only if no attempt newer than the retry has been |
| 117 | + stored, because such an attempt means the message failed again afterwards |
| 118 | +- an attempt moves the status only if it is strictly newer than the newest attempt already stored, |
| 119 | + so a redelivery of the attempt already there cannot undo a resolve or an archive |
| 120 | + |
| 121 | +## Where things live |
| 122 | + |
| 123 | +| | | |
| 124 | +| --- | --- | |
| 125 | +| `ServiceControl.Infrastructure/Ingestion/IngestionPipeline.cs` | The channels, the assembler and the writers | |
| 126 | +| `ServiceControl.Infrastructure/Ingestion/IngestionSettingsReader.cs` | Reading, validating and resolving the settings above | |
| 127 | +| `ServiceControl.Infrastructure/Ingestion/Metrics/` | The metric scopes both instances report through | |
| 128 | +| `ServiceControl/Operations/ErrorIngestion.cs` | Transport, watchdog and fault policy for the error queue | |
| 129 | +| `ServiceControl.Audit/Auditing/AuditIngestion.cs` | The same for the audit queue | |
0 commit comments