Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
53 commits
Select commit Hold shift + click to select a range
fd6b5d4
feat: indexing pipeline drain
Totodore Sep 4, 2026
adffac8
feat: nats source
Totodore Sep 4, 2026
137846c
feat: indexing pipeline drain
Totodore Sep 4, 2026
21c16d8
Merge branch 'feat-pipeline-drain' into feat-durable-nats-source
Totodore Sep 4, 2026
ef62d74
feat: indexing pipeline drain
Totodore Sep 4, 2026
d2790f1
Merge branch 'feat-pipeline-drain' into feat-durable-nats-source
Totodore Sep 4, 2026
3b17597
feat: indexing pipeline drain
Totodore Sep 4, 2026
430bde4
feat: indexing pipeline drain
Totodore Sep 4, 2026
6b0adc1
feat: indexing pipeline drain
Totodore Sep 4, 2026
2ef145e
feat: indexing pipeline drain
Totodore Sep 4, 2026
88d1242
feat: indexing pipeline drain
Totodore Sep 4, 2026
b4b2b60
feat: indexing pipeline drain
Totodore Sep 4, 2026
8e311e7
feat: indexing pipeline drain
Totodore Sep 4, 2026
c007b10
feat: indexing pipeline drain
Totodore Sep 4, 2026
d6588f3
feat: indexing pipeline drain
Totodore Sep 4, 2026
946b518
feat: indexing pipeline drain
Totodore Sep 4, 2026
b615e9f
feat: indexing pipeline drain
Totodore Sep 4, 2026
cd9b6cc
feat: indexing pipeline drain
Totodore Sep 4, 2026
a23df94
fix(indexers): disable new pipeline spawning when draining
Totodore Sep 7, 2026
90fd1a1
fix(indexers): drive the drain supervise loop at 1s rather than 30s)
Totodore Sep 7, 2026
94a28c2
fix(doc): rephrase the indexer shutdown drain timeout description
Totodore Sep 7, 2026
169e7a8
fix(doc): wait for DrainPipeline so that non-draining pipeline are ki…
Totodore Sep 7, 2026
1574496
Merge branch 'main' into feat-pipeline-drain
Totodore Sep 7, 2026
f77e519
fix(doc): wait for DrainPipeline so that non-draining pipeline are ki…
Totodore Sep 7, 2026
69d4f64
Merge branch 'feat-pipeline-drain' into feat-durable-nats-source
Totodore Sep 7, 2026
39a89c8
fix: config
Totodore Sep 7, 2026
3b5e971
fix: config
Totodore Sep 7, 2026
2d9c42c
feat: add tests
Totodore Sep 7, 2026
8ac7c8d
feat: add tests
Totodore Sep 7, 2026
928f4a9
fix: fmt
Totodore Sep 7, 2026
a19e18f
fix: nats
Totodore Sep 8, 2026
9453702
fix: nat source
Totodore Sep 8, 2026
17f226c
feat: nats source
Totodore Sep 8, 2026
67839a1
fix: comments
Totodore Sep 8, 2026
96efbf1
fix: comments
Totodore Sep 8, 2026
fcfcca9
fix: remove the ack policy all
Totodore Sep 8, 2026
55b5e53
fix: remove the ack policy all
Totodore Sep 8, 2026
3804790
feat: nats source
Totodore Sep 9, 2026
2667e4c
feat: nats docs
Totodore Sep 9, 2026
4fddd3a
fix: source config
Totodore Sep 9, 2026
29621c9
feat(indexers): fixes
Totodore Sep 9, 2026
76ccf71
feat(indexers): remove useless hb
Totodore Sep 9, 2026
857a3fe
fix: minor fixes
Totodore Sep 9, 2026
3be262c
fix: minor fixes
Totodore Sep 9, 2026
f2d60fb
feat: improve tracing
Totodore Sep 9, 2026
2ed8834
Merge branch 'main' into feat-pipeline-drain
Totodore Sep 24, 2026
25bcbdc
Merge remote-tracking branch 'origin/main' into feat-pipeline-drain
Totodore Sep 29, 2026
b414378
Merge branch 'feat-pipeline-drain' into feat-durable-nats-source
Totodore Sep 29, 2026
65f6395
Merge branch 'main' into feat-pipeline-drain
Totodore Sep 29, 2026
b55659b
Merge branch 'main' into feat-pipeline-drain
Totodore Oct 5, 2026
25900df
Merge branch 'main' into feat-pipeline-drain
Totodore Oct 7, 2026
c524d47
Merge remote-tracking branch 'origin/main' into feat-durable-nats-source
Totodore Oct 7, 2026
8e24d3a
Merge remote-tracking branch 'origin/feat-pipeline-drain' into feat-d…
Totodore Oct 8, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ MAP_HOST_POSTGRES=0.0.0.0
MAP_HOST_PULSAR=0.0.0.0
MAP_HOST_KAFKA=0.0.0.0
MAP_HOST_ZOOKEEPER=0.0.0.0
MAP_HOST_NATS=0.0.0.0
MAP_HOST_AZURITE=0.0.0.0
MAP_HOST_GRAFANA=0.0.0.0
MAP_HOST_JAEGER=0.0.0.0
Expand Down
3 changes: 3 additions & 0 deletions .github/workflows/coverage.yml
Original file line number Diff line number Diff line change
Expand Up @@ -151,6 +151,9 @@ jobs:
- name: Run Pulsar service
run: DOCKER_SERVICES=pulsar make docker-compose-up

- name: Run NATS service
run: DOCKER_SERVICES=nats make docker-compose-up

- name: Install Rust
run: rustup toolchain install
working-directory: ./quickwit
Expand Down
5 changes: 5 additions & 0 deletions LICENSE-3rdparty.csv
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@ async-channel,https://github.com/smol-rs/async-channel,Apache-2.0 OR MIT,Stjepan
async-compression,https://github.com/Nullus157/async-compression,MIT OR Apache-2.0,"Wim Looman <wim@nemo157.com>, Allen Bui <fairingrey@gmail.com>"
async-io,https://github.com/smol-rs/async-io,Apache-2.0 OR MIT,Stjepan Glavina <stjepang@gmail.com>
async-lock,https://github.com/smol-rs/async-lock,Apache-2.0 OR MIT,Stjepan Glavina <stjepang@gmail.com>
async-nats,https://github.com/nats-io/nats.rs,Apache-2.0,"Tomasz Pietrek <tomasz@nats.io>, Casper Beyer <caspervonb@pm.me>"
async-process,https://github.com/smol-rs/async-process,Apache-2.0 OR MIT,Stjepan Glavina <stjepang@gmail.com>
async-signal,https://github.com/smol-rs/async-signal,Apache-2.0 OR MIT,John Nunley <dev@notgull.net>
async-speed-limit,https://github.com/tikv/async-speed-limit,MIT OR Apache-2.0,The TiKV Project Developers
Expand Down Expand Up @@ -674,9 +675,11 @@ serde_core,https://github.com/serde-rs/serde,MIT OR Apache-2.0,"Erick Tryzelaar
serde_derive,https://github.com/serde-rs/serde,MIT OR Apache-2.0,"Erick Tryzelaar <erick.tryzelaar@gmail.com>, David Tolnay <dtolnay@gmail.com>"
serde_json,https://github.com/serde-rs/json,MIT OR Apache-2.0,"Erick Tryzelaar <erick.tryzelaar@gmail.com>, David Tolnay <dtolnay@gmail.com>"
serde_json_borrow,https://github.com/PSeitz/serde_json_borrow,MIT,Pascal Seitz <pascal.seitz@gmail.com>
serde_nanos,https://github.com/caspervonb/serde_nanos,MIT OR Apache-2.0,Casper Beyer <caspervonb@pm.me>
serde_path_to_error,https://github.com/dtolnay/path-to-error,MIT OR Apache-2.0,David Tolnay <dtolnay@gmail.com>
serde_plain,https://github.com/mitsuhiko/serde-plain,MIT OR Apache-2.0,Armin Ronacher <armin.ronacher@active-4.com>
serde_qs,https://github.com/samscott89/serde_qs,MIT OR Apache-2.0,Sam Scott <sam@osohq.com>
serde_repr,https://github.com/dtolnay/serde-repr,MIT OR Apache-2.0,David Tolnay <dtolnay@gmail.com>
serde_spanned,https://github.com/toml-rs/toml,MIT OR Apache-2.0,The serde_spanned Authors
serde_urlencoded,https://github.com/nox/serde_urlencoded,MIT OR Apache-2.0,Anthony Ramine <n.oxyde@gmail.com>
serde_with,https://github.com/jonasbb/serde_with,MIT OR Apache-2.0,"Jonas Bushart, Marcin Kaźmierczak"
Expand Down Expand Up @@ -776,6 +779,7 @@ tokio-retry2,https://github.com/naomijub/tokio-retry,MIT,"Julia Naomi <jnboeira@
tokio-rustls,https://github.com/rustls/tokio-rustls,MIT OR Apache-2.0,The tokio-rustls Authors
tokio-stream,https://github.com/tokio-rs/tokio,MIT,Tokio Contributors <team@tokio.rs>
tokio-util,https://github.com/tokio-rs/tokio,MIT,Tokio Contributors <team@tokio.rs>
tokio-websockets,https://github.com/Gelbpunkt/tokio-websockets,MIT,The tokio-websockets Authors
toml,https://github.com/toml-rs/toml,MIT OR Apache-2.0,The toml Authors
toml_datetime,https://github.com/toml-rs/toml,MIT OR Apache-2.0,The toml_datetime Authors
toml_edit,https://github.com/toml-rs/toml,MIT OR Apache-2.0,The toml_edit Authors
Expand All @@ -801,6 +805,7 @@ tracing-serde,https://github.com/tokio-rs/tracing,MIT,Tokio Contributors <team@t
tracing-subscriber,https://github.com/tokio-rs/tracing,MIT,"Eliza Weisman <eliza@buoyant.io>, David Barsky <me@davidbarsky.com>, Tokio Contributors <team@tokio.rs>"
triomphe,https://github.com/Manishearth/triomphe,MIT OR Apache-2.0,"Manish Goregaokar <manishsmail@gmail.com>, The Servo Project Developers"
try-lock,https://github.com/seanmonstar/try-lock,MIT,Sean McArthur <sean@seanmonstar.com>
tryhard,https://github.com/EmbarkStudios/tryhard,MIT OR Apache-2.0,Embark <opensource@embark-studios.com>
ttl_cache,https://github.com/stusmall/ttl_cache,MIT OR Apache-2.0,Stu Small <stuart.alan.small@gmail.com>
twox-hash,https://github.com/shepmaster/twox-hash,MIT,Jake Goulding <jake.goulding@gmail.com>
typeid,https://github.com/dtolnay/typeid,MIT OR Apache-2.0,David Tolnay <dtolnay@gmail.com>
Expand Down
11 changes: 11 additions & 0 deletions config/tutorials/stackoverflow/nats-source.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
#
# NATS source config file.
#
version: 0.8
source_id: nats-source
source_type: nats
params:
uris:
- nats://localhost:4222
stream: stackoverflow
consumer: quickwit-consumer
19 changes: 19 additions & 0 deletions config/tutorials/stackoverflow/send_messages_to_nats.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,19 @@
import asyncio

import nats


async def main():
client = await nats.connect("nats://localhost:4222")
jetstream = client.jetstream()

with open("stackoverflow.posts.transformed-10000.json", encoding="utf8") as file:
for i, line in enumerate(file):
await jetstream.publish("stackoverflow.posts", line.strip().encode("utf-8"))
if i % 1000 == 0:
print(f"{i}/10000 messages sent.")

await client.close()


asyncio.run(main())
15 changes: 15 additions & 0 deletions docker-compose.yml
Original file line number Diff line number Diff line change
Expand Up @@ -127,6 +127,21 @@ services:
timeout: 10s
retries: 100

nats:
image: nats:${NATS_VERSION:-2.10-alpine}
container_name: nats
command: --jetstream --http_port 8222
ports:
- "${MAP_HOST_NATS:-127.0.0.1}:4222:4222"
profiles:
- all
- nats
healthcheck:
test: ["CMD", "wget", "-q", "-O-", "http://localhost:8222/healthz"]
interval: 1s
timeout: 5s
retries: 100

azurite:
image: mcr.microsoft.com/azure-storage/azurite:${AZURITE_VERSION:-3.24.0}
container_name: azurite
Expand Down
1 change: 1 addition & 0 deletions docs/configuration/node-config.md
Original file line number Diff line number Diff line change
Expand Up @@ -247,6 +247,7 @@ This section contains the configuration options for an indexer. The split store
| `enable_otlp_endpoint` | If true, enables the OpenTelemetry exporter endpoint to ingest logs and traces via the OpenTelemetry Protocol (OTLP). | `false` |
| `cpu_capacity` | Advisory parameter used by the control plane. The value can expressed be in threads (e.g. `2`) or in term of millicpus (`2000m`). The control plane will attempt to schedule indexing pipelines on the different nodes proportionally to the cpu capacity advertised by the indexer. It is NOT used as a limit. All pipelines will be scheduled regardless of whether the cluster has sufficient capacity or not. The control plane does not attempt to spread the work equally when the load is well below the `cpu_capacity`. Users who need a balanced load on all of their indexer nodes can set the `cpu_capacity` to an arbitrarily low value as long as they keep it proportional to the number of threads available. | `num threads available` |
| `enable_cooperative_indexing` | Enable sharing resources more efficiently when the number of indexes actively written to is significantly higher than the number of cores but might decrease the overall indexing throughput. | `false` |
| `shutdown_drain_timeout` | Time budget granted to each indexing pipeline to drain gracefully before its remaining actors are killed, on node shutdown and when the control plane tears down a pipeline. Set it above the largest `commit_timeout_secs` plus the time to upload and publish the final splits. Beware that on shutdown, draining only starts after the ingester and compactor decommissions complete, so the deployment's shutdown grace period (e.g. `terminationGracePeriodSeconds` on Kubernetes) must exceed the largest decommission timeout *plus* this value — about 600 seconds with the default timeouts. Can be overridden with the `QW_INDEXER_SHUTDOWN_DRAIN_TIMEOUT` environment variable. | `300s` |

Set the `QW_INDEXING_MAX_WRITE_THROUGHPUT` environment variable to limit the aggregate indexing IO throughput on a node. It accepts human-readable byte sizes per second, such as `500mb`, and is unlimited by default.

Expand Down
103 changes: 101 additions & 2 deletions docs/configuration/source-config.md
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@ The source ID is a string that uniquely identifies the source within an index. I

## Source type

The source type designates the kind of source being configured. As of version 0.5, available source types are `ingest-api`, `kafka`, `kinesis`, and `pulsar`. The `file` type is also supported but only for local ingestion from [the CLI](/docs/reference/cli.md#tool-local-ingest).
The source type designates the kind of source being configured. Available source types are `ingest-api`, `kafka`, `kinesis`, `nats`, `pubsub`, and `pulsar`. The `file` type is also supported but only for local ingestion from [the CLI](/docs/reference/cli.md#tool-local-ingest).

## Source parameters

Expand Down Expand Up @@ -183,6 +183,105 @@ EOF
quickwit source create --index my-index --source-config source-config.yaml
```

### NATS source

A NATS source reads data from a [NATS JetStream](https://docs.nats.io/nats-concepts/jetstream) stream through a durable consumer. Each message carries one payload in the source's [input format](#input-format): a single JSON object (`json`, the default), a plain text document (`plain_text`), or an OTLP export request whose log records or spans are each indexed as a separate document (`otlp_*` formats). Payloads must not exceed 1 MiB.

A tutorial is available [here](/docs/ingest-data/nats.md).

The durable consumer is provisioned externally so its lifecycle, subject filters, deliver policy, and ack tuning belong to whoever provisioned it.

Delivery is **exactly-once on planned teardowns and at-least-once on crashes**. On a planned teardown (node shutdown, pipeline reassignment on a `num_pipelines` change), the pipeline is drained first: the source stops pulling, the in-flight messages are committed, published, and acknowledged before the pipeline stops, so nothing is indexed twice. Messages the pipeline had prefetched but not processed are negatively acknowledged so the remaining pipelines pick them up immediately. The drain runs under a time budget, the indexer's [`shutdown_drain_timeout`](node-config.md#indexer-configuration). A drain that cannot finish within it (e.g. the object storage or the metastore is unavailable) is abandoned and delivery degrades to at-least-once, as on a crash. After a crash, the unacknowledged messages are redelivered after `ack_wait` and indexed again, as duplicates.

The source never waits on unreachable NATS servers: acknowledgments and the negative acknowledgments sent at teardown are bounded by a 10 s timeout, unconfirmed acknowledgments stay pending and are retried at the next split publication, and a pipeline can still be stopped or reassigned during an outage.

#### Consumer invariants

These are properties of the consumer, not of the source. Only the ack policy is enforced when the source is created. The pull limits and `max_deliver` are checked when a pipeline starts and only produce a warning in the indexer logs. The rest are not checked at all, and getting them wrong looks like Quickwit being slow or duplicating rather than like a consumer misconfiguration.

**`ack_policy` must be `explicit`.** The source acknowledges each message individually once the split containing it is published, and waits for the server to confirm the acknowledgment — the confirmation is what tells the drain that the pipeline is empty. The consumer's ack floor is therefore the resume point, and it is the only progress state that matters. Any other policy is rejected when the source is created.

**`ack_wait` must exceed the end-to-end publish latency.** The timer starts at delivery and has to outlast three terms:

```
ack_wait > commit_timeout + split upload and publish + ack round trip
```

If `ack_wait` is shorter, NATS redelivers messages that are still being indexed; they are then indexed twice. 5 minutes with a 60 s `commit_timeout` leaves a wide margin.

It is also the *recovery* time from a lost delivery, so it should not be arbitrarily large either. If a message is delivered but never arrives (a dropped connection, a slow-consumer disconnect) nothing brings it back until the timer expires and the pipeline simply idles.

**`max_ack_pending` should be `-1`.** Nothing is acknowledged until a split is published, so a whole commit window is always ack-pending. This setting is a hard throughput cap rather than a safety valve:

```
achievable rate ≈ max_ack_pending / commit_timeout documents per second
```

Measured at four pipelines with a 10 s commit timeout: `1000` gave 97 documents per second where the formula predicts 100, and `20000` gave 2,115 where it predicts 2,000. A bound that looks generous for an ordinary queue consumer throttles indexing. If the in-flight window has to be bounded, size it above `throughput × commit_timeout` rather than at a small absolute number, and note the trade: unlimited also makes the crash-replay window the whole delivered span above the ack floor rather than a bounded slice.

**`max_deliver` should stay unlimited (`-1`).** Redelivery is what recovers from a crash or a lost delivery. A finite value turns it into data loss: a message that exhausts it is never delivered again and is silently skipped.

**`max_batch`, `max_bytes` must not be below what the source pulls.** Each pull request asks for up to `QW_NATS_PULL_MAX_MESSAGES_PER_BATCH` messages and `QW_NATS_PULL_MAX_BYTES_PER_BATCH` bytes. A consumer limit below any of those makes the server reject every pull request and the pipeline idles. The source warns at start.

**The consumer must outlive the source.** Deleting it, or letting an `inactive_threshold` expire, while a source is bound to it does not stop the pipelines: the server answers their pull requests with "no responders", which is also what a JetStream outage produces, so the source retries indefinitely and warns. Recreate the consumer under the same name to resume, or disable the source.

**`deliver_policy`** is usually `all`, so a new consumer indexes the stream from the start.

**Give each consumer a single `filter_subject`.** A consumer configured with several filters through NATS 2.10's multi-filter `filter_subjects` appears to take a much slower path in the broker: on a stream whose subjects are interleaved, sixteen multi-filter consumers took 7.5 times longer than the same sixteen consumers with one filter each, with the broker saturated and the indexers idle.

#### What can be tuned on the Quickwit side

| Setting | Where | Effect |
| --- | --- | --- |
| `num_pipelines` | source config | Pipelines bound to the consumer. Scales cleanly on large messages: 52 → 147 MiB/s from one to four pipelines at 512 KiB. On small messages (1 KiB) the acknowledgment path caps the aggregate, and it is worth only about 1.15× from one to sixteen — there, parallelism has to come from more consumers. |
| [`commit_timeout_secs`](index-config.md#indexing-settings) | index config | Sets the size of the always-ack-pending window, so it interacts with `max_ack_pending` and with `ack_wait`. |
| `QW_NATS_PULL_MAX_BYTES_PER_BATCH` | env, default 10 MiB | Bounds the bytes the server may push per pull request. It must stay well under the server's per-connection `max_pending` (64 MiB by default), or the server declares a slow consumer and closes the connection; the messages are already counted as delivered, so the pipeline then idles for a full `ack_wait`. It must also stay above the server's `max_payload`, or larger messages are never delivered, and at or below the consumer's `max_bytes`; the source warns at start otherwise. 10 MiB and 20 MiB measure identically. |
| `QW_NATS_PULL_MAX_MESSAGES_PER_BATCH` | env, default 100,000 | Messages per pull request. Can be used along `QW_NATS_PULL_MAX_BYTES_PER_BATCH` to limit a batch if there is too much contention around acknowledgements. The client buffers eight such batches per subscription; messages beyond that are dropped client-side and come back after `ack_wait`. |

#### Scaling

Several pipelines can bind to the same consumer and NATS load-balances the messages across them, so scaling is a plain `num_pipelines` update and the control plane places the pipelines across the indexers of the cluster.

#### Acknowledgment cost

One confirmed acknowledgment per message is the source's dominant broker cost, and it is what decides whether the source or the indexer is the bottleneck.

- At **512 KiB** per message it is invisible. The NATS and Kafka sources measured within a percent of each other across the whole pipeline sweep.
- At **1 KiB** it is the ceiling. Sixteen per-tenant consumers reached 72 MiB/s while an equivalent Kafka source reached 150, and the broker spent 5.1 of the host's cores on acknowledgments against Quickwit's 4.7 on indexing — about ten times the broker CPU per document that a periodic offset commit costs. Priced on the broker alone, acknowledgment divides throughput by 6.5 at this size.

#### Monitoring

Being durable, the consumer is observable through NATS's own monitoring (`nats consumer info`, exporters): `num_pending` is the indexing lag and `num_ack_pending` the in-flight window, both available even while the pipelines are down. The source also reports `num_pending_acks` (messages indexed but whose split is not published yet) in its observable state.

When a message carries a W3C `traceparent` header and the indexers export their traces (see [distributed tracing](../distributed-tracing/plug-quickwit-to-jaeger.md)), the publisher's trace extends to Quickwit: a `process_nats_message` span covers the message's processing until it is acknowledged.

**NATS source parameters**

| Property | Description | Default value |
| --- | --- | --- |
| `uris` | List of NATS server URIs (e.g. `nats://localhost:4222`). | required |
| `stream` | Name of the JetStream stream to consume. | required |
| `consumer` | Name of the pre-provisioned durable consumer to bind to. | required |
| `tls` | TLS options: `ca_certificates_path` (PEM file whose root certificates are trusted instead of the system ones), and `client_certificate_path` + `client_key_path` (PEM files, set together) for mutual TLS. TLS itself is enabled by connecting to `tls://` URIs. The files are read when a connection is established: by the indexer nodes running the source, and by the node serving the source creation or update, which checks connectivity. | optional |
| `authentication` | Authentication parameters: either `user_password` (with `user` and `password`) or `token`. | optional |

*Adding a NATS source to an index with the [CLI](../reference/cli.md#source)*

```bash
cat << EOF > source-config.yaml
version: 0.8
source_id: my-nats-source
source_type: nats
num_pipelines: 2
params:
uris:
- nats://localhost:4222
stream: my-stream
consumer: my-consumer
EOF
./quickwit source create --index my-index --source-config source-config.yaml
```

### Pulsar source

A Puslar source reads data from one or several Pulsar topics. Each message in topic(s) must hold a JSON object.
Expand Down Expand Up @@ -216,7 +315,7 @@ EOF

## Number of pipelines

The `num_pipelines` parameter is only available for distributed sources like Kafka, GCP PubSub, and Pulsar.
The `num_pipelines` parameter is only available for distributed sources like Kafka, GCP PubSub, NATS, and Pulsar.

It defines the number of pipelines to run on a cluster for the source. The actual placement of these pipelines on the different indexer
will be decided by the control plane.
Expand Down
Loading