From 78ba113c04fcbd5148da23ed4149062f65cd1f6c Mon Sep 17 00:00:00 2001 From: Zdravko Donev Date: Mon, 14 Sep 2026 02:21:35 +0300 Subject: [PATCH 1/8] RDSC-6038: Add RDI 2.0.0 release notes --- config.toml | 2 +- .../release-notes/rdi-2-0-0.md | 78 +++++++++++++++++++ 2 files changed, 79 insertions(+), 1 deletion(-) create mode 100644 content/integrate/redis-data-integration/release-notes/rdi-2-0-0.md diff --git a/config.toml b/config.toml index 230dee5230..d4365ed220 100644 --- a/config.toml +++ b/config.toml @@ -97,7 +97,7 @@ rdi_redis_gears_version = "1.2.6" rdi_debezium_server_version = "2.3.0.Final" rdi_db_types = "cassandra|mysql|oracle|postgresql|sqlserver" rdi_cli_latest = "latest" -rdi_current_version = "1.19.1" +rdi_current_version = "2.0.0" [params.clientsConfig] diff --git a/content/integrate/redis-data-integration/release-notes/rdi-2-0-0.md b/content/integrate/redis-data-integration/release-notes/rdi-2-0-0.md new file mode 100644 index 0000000000..a4a6ce3f40 --- /dev/null +++ b/content/integrate/redis-data-integration/release-notes/rdi-2-0-0.md @@ -0,0 +1,78 @@ +--- +Title: Redis Data Integration release notes 2.0.0 (September 2026) +alwaysopen: false +categories: +- docs +- operate +- rs +description: | + Multiple sources in one pipeline, Flink as the default processor, and simpler source mTLS configuration. + Database-scoped secrets, per-source reset and cleanup, expanded API diagnostics, and security updates. +linkTitle: 2.0.0 (September 2026) +toc: 'true' +weight: 967 +--- + +## What's New in 2.0.0 + +### Breaking Changes + +- **API v1 pipeline actions apply immediately instead of queuing**: + - `POST /pipelines`, `PATCH /pipelines`, `/pipelines/undeploy`, `/pipelines/sources/*`, `/pipelines/targets/*`, `/pipelines/processors/*`, `/pipelines/secret-providers/*`, `/pipelines/start`, `/pipelines/stop`, `/pipelines/reset` now update the Pipeline resource directly instead of posting to an internal task queue. Kubernetes errors are reported with their original status codes, rather than a generic 500. + - A request now applies its changes immediately rather than waiting behind earlier in-flight requests, so the most recent request wins if several are issued in quick succession without waiting for each to complete. Polling `GET /actions/{action_id}` for an action id that a later request has superseded now returns an unknown action error instead of that action's own status. + - `POST /pipelines` and `PATCH /pipelines` now also validate target connectivity similarly to their API v2 counterparts, and surface collector API failures as a `422`, `502`, `503`, or `504` error response. A request that previously succeeded despite an unreachable target now fails with a validation error. +- **API v1 trace endpoint returns Not Implemented**: `POST /trace/start` now immediately returns `501 Not Implemented` instead of accepting the request and queuing a trace that would never run. +- **API v2 pipeline responses omit optional fields that have no value**: The pipeline, pipeline status, create, update, patch, start, stop, and reset responses no longer carry an optional field whose value is null; the field is left out instead. This affects `status_changed_at` and an error's `remediation`. The OpenAPI schema is unchanged, since neither field was ever required, so clients generated from it are unaffected; a client that reads either key directly has to treat it as absent rather than null. +- **A source can no longer set its own `topic.prefix`**: The `topic.prefix` property is now rejected in a source's `advanced.source` section, since RDI derives the topic prefix from the source name. Remove the property from your configuration and set each job's `server_name` to the derived prefix: the source name, `rdi` for the source of an upgraded single-source pipeline handled by the Debezium collector, and the Spanner instance ID for the source of such an upgraded pipeline handled by the Flink collector. +- **API v2 metric collections name data streams after their table**: The `data_streams.streams` keys of a metric collection are now the source-qualified table name (`mysql.inventory.addresses`), matching what the dead-letter queue endpoints return, instead of the Redis stream name (`{rdi}:inventory.addresses`). A client that keys off the previous form has to be updated. API v1 statistics are unchanged. +- **Cassandra is no longer a source database type**: `cassandra` has been removed from the source database types that `redis-di scaffold`, the configuration template endpoint, and `rdi-admin install` offer, and a source whose connection type is `cassandra` no longer passes validation. RDI never actually supported Cassandra as source, since the Debezium Cassandra connector has to run on each Cassandra node. +- **The Flink processor is the default processor**: A pipeline whose `processors` section does not set `type` now deploys the Flink processor instead of the classic one. Set `processors.type` to `classic` to keep deploying the classic processor. + +### New Features + +- **Multiple sources in one pipeline**: A pipeline can ingest data from several source databases, of the same or different types, into one Redis target. Each source has its own name, connection settings, credentials, and collector. Use API v2 or the `redis-di` CLI to manage multi-source pipelines; API v1 supports only single-source pipelines. + +- **Transformation jobs can select several tables (Flink processor)**: A job's `server_name`, `db`, `schema`, and `table` source matchers now accept a list of values in addition to a single one, and an entry prefixed with `regex:` is matched as an anchored regular expression, so one job can handle many tables. An entry without the prefix is matched literally, so existing jobs are unaffected, including table names containing regular expression characters. Two jobs whose matchers select the same table are rejected. Only the Flink processor supports this syntax, so a job using it is rejected for the classic processor. +- **Database-scoped pipeline secret keys**: Pipeline secrets in API v2 and the `redis-di` CLI now use database-independent keys (`USERNAME`, `PASSWORD`, `CACERT`, `CERT`, `KEY`, `KEY_PASSWORD`) together with a `db` parameter (`--db` in the CLI) that names the database the secret belongs to: a source name, or `target`. The previous scope-prefixed keys (`SOURCE_DB_*`, `TARGET_DB_*`) remain accepted for single-source and target secrets, used without the `db` parameter. +- **Richer validation errors from API v1 pipeline configuration endpoints**: A `422` response from `/pipelines/undeploy`, `/pipelines/sources/*`, `/pipelines/targets/*`, `/pipelines/processors/*`, and `/pipelines/secret-providers/*` may now include an `errors` array with structured per-field detail, in addition to the existing `detail` message, matching the shape already returned by `POST /pipelines` and `PATCH /pipelines`. Existing clients that only read `detail` are unaffected. +- **Pipeline components report the source they are associated with**: Each collector component in an API v2 pipeline or pipeline status response now has a `source` field naming the pipeline source it is associated with. It is empty for components that are not per-source, such as the processor and the Collector API. +- **A follower installation keeps its pipelines instead of losing them on a leadership handover**: In a high-availability RDI deployment with leader election enabled, the installation that is not currently active now keeps its Pipeline resources and reports them with a new `standby` status, instead of deleting them and recreating them on the next handover. Mutating a pipeline (create, update, patch, delete, start, stop, reset, or start/stop a source) on a standby installation now returns `503 Service Unavailable`; reads keep working. +- **Scaffolded configurations use a descriptive source name**: `redis-di scaffold` and the configuration template endpoint now name the generated source after the database type (for example `mysql`) and reference the matching source-prefixed secrets (for example `${MYSQL_DB_USERNAME}`), instead of the generic `source` name with `${SOURCE_DB_*}` references, so the config, the injected environment variables, and the secrets set with `--db ` always line up. A custom name can be passed with the new `--source-name` flag (`source_name` query parameter). The silent `rdi-admin install` writes its source secrets under the same derived names and accepts an optional `sources.default.name` key; when it reuses an existing configuration instead of scaffolding one, the secrets keep the legacy `source` name unless `sources.default.name` says otherwise. It also stores source and target certificates under their canonical file names (`ca.crt`, `client.crt`, `client.key`) in the `-db-ssl` secret, matching what the API writes, instead of the names of the files they were uploaded from. +- **TypeScript SDK published to npmjs**: The RDI API TypeScript SDK is now publicly available on npmjs as `@rdi-ui/sdk` under the MIT license, so it no longer requires access to the internal GitHub Packages registry: `npm add @rdi-ui/sdk`. +- **More collector diagnostics in API v2 metric collections**: The `GET /pipelines/{name}/metric-collections` endpoint now reports the full set of numeric Debezium collector metrics per source: event counts by operation, filtered and erroneous events, queue capacity and byte usage, the last processed transaction id, MySQL and MariaDB binlog health counters and GTID set, MongoDB primary elections, and Oracle LogMiner SCN, lag, performance, and error metrics. +- **Pipeline components report their externally-reachable endpoints**: Each component in an API v2 pipeline or pipeline status response now has an `external_endpoints` field listing the URLs at which that component is reachable from outside the cluster, discovered from its Ingress resources, in addition to the existing internal `metrics_endpoints`. +- **Source mTLS without Debezium keystore settings**: A MySQL, MariaDB, or MongoDB pipeline now presents the source client certificate without setting `database.ssl.keystore`/`mongodb.ssl.keystore` and the matching password in `advanced.source`. The collector points the connector at the keystore RDI builds from the source's certificate secrets, whether stored as pipeline secrets or supplied by a secret provider, and explicit `advanced.source` settings keep overriding the derived values. +- **Reset a single source of a multi-source pipeline**: `POST /pipelines/{name}/reset` accepts an optional `source` query parameter and `redis-di reset` an optional `--source` flag, naming one existing source that is not of type `external`. The reset then deletes only that source's data, including its change streams, offsets, schema history, dead letter queue entries, statistics, deduplication state, and record counters, and leaves every other source's data intact. The whole pipeline stops while the reset runs and starts again afterwards, as it already does for full pipeline reset. +- **Removing a source from a pipeline deletes its internal RDI data**: Removing a source with `PUT` or `PATCH /pipelines/{name}` now deletes the internal RDI data that source leaves behind, including its change streams, offsets, schema history, dead letter queue entries, statistics, deduplication state, and record counters. That data used to be kept indefinitely, and only resetting the whole pipeline removed it. The other sources keep their data, and the whole pipeline stops while the removed source's data is deleted and starts again afterwards. Records already written to the target Redis database are retained. + +### Bug Fixes + +- **RDI database client certificates isolated from source certificates**: The Debezium collector now presents its RDI database client certificate from a dedicated keystore instead of the shared default keystore, so the source database and RDI database client identities can no longer interfere with each other under mTLS. +- **Nested processor properties shown incorrectly in `redis-di describe`**: Nested objects and arrays in a processor's advanced properties are now shown as JSON instead of a raw internal format, so the `redis-di describe` output is valid and easy to read. +- **Bulk insert right after a snapshot could restart the Flink processor (Flink processor)**: A bulk insert arriving just after a snapshot completed could crash the Flink job with `Have records for a split that was not registered` while the source transitioned from snapshot to CDC processing. The job recovered automatically, but the pending change records were delayed until the restarted job claimed them. The source reader now skips record batches of splits that finished during the transition instead of failing on them. +- **Stream could stop being processed after a Flink processor restart (Flink processor)**: A narrow timing window around a checkpoint could leave a stream permanently unassigned to any reader after a restart (for example following a TaskManager failure), silently halting ingestion for that stream until the pipeline was redeployed. The processor now confirms with every reader before treating a stream as no longer needed. +- **Unreachable source database reported as a validation error**: When RDI cannot list the source tables during pipeline validation, for example because the source database refuses the connection, the failure is now reported as a `422` validation error in the `errors` array, instead of a `422`, `502`, or `504` response carrying only a `detail` message. An unavailable collector API is still reported as a `503` response so the request can be retried. +- **No configuration template for the Redis database type**: `GET /pipelines/config/templates/ingest/redis` now returns the documented `501 Not Implemented`, instead of a `200` response carrying a configuration that named `redis` as the source connection type and left the source port empty. Redis cannot be used as a pipeline source. +- **A database flavor that does not match the database type is rejected**: `GET /pipelines/config/templates/ingest/{db_type}` now returns a `422` response when `db_flavor` names a MongoDB flavor and `db_type` is not `mongodb`, instead of a `200` response carrying a configuration that ignored the flavor. + +### Improvements + +- **API response body logging**: The RDI API now logs response bodies alongside request payloads. Bodies are logged at DEBUG level, while request and response metadata stays at INFO, and large bodies are truncated to keep logs manageable. +- **Documented scaffolded configurations**: The configuration that `redis-di scaffold`, the configuration template endpoint, and `rdi-admin install` generate now documents every property it offers, with a description above each one, links to the pipeline configuration documentation, and the default value where one exists. It also covers Snowflake and it no longer suggests properties that the source database does not support. The target port now defaults to `6379` rather than a placeholder, so the generated configuration passes schema validation before any property is filled in. The per-database example configurations that `rdi-admin install` wrote alongside it (`config.yaml.mysql.example` and its siblings) have been removed, since the generated configuration and the documentation now cover the same ground. +- **Documented transformation job template**: The job template that `GET /pipelines/jobs/templates/ingest` returns now documents the properties of the source, transform, and output blocks, with a description above each one and links to the transformation documentation, instead of naming a handful of them in a bare skeleton. It is a valid job as returned, with only the mandatory properties left uncommented. The endpoint accepts a new optional `source_name` query parameter that names the source the job reads from, matching the parameter of the configuration template endpoint. `rdi-admin install` now writes the same template to `jobs/job.yaml`, named after the source it scaffolds, in place of the example jobs it used to copy there. +- **Debezium collector processes captured records on four threads**: The Debezium collector now sets `record.processing.threads` to 4 rather than letting Debezium size the pool by itself, which raises snapshot throughput considerably. Lower it through a source's `advanced.source.record.processing.threads` when fewer CPUs are available to the collector. +- **Faster pipeline reset**: Resetting a pipeline is faster and no longer creates a Kubernetes Job. + +### Security + +- **Hardened API log redaction**: Ensured secrets, credentials, tokens, and connection strings are masked in the request and response bodies the RDI API logs, with login and pipeline secret payloads fully redacted. +- **Flink processor and collector security updates**: Resolved `CVE-2024-57699`, `CVE-2025-12183`, `CVE-2025-27820`, `CVE-2025-55163`, `CVE-2025-66566`, `CVE-2025-68973`, `CVE-2026-35194`, `CVE-2026-42198`, `CVE-2026-42583`, `CVE-2026-44249`, `CVE-2026-45416`, `CVE-2026-45447`, `CVE-2026-50010`, `CVE-2026-54291`, `CVE-2026-54512`, `CVE-2026-54513`, `CVE-2026-59901`, and `GHSA-r7wm-3cxj-wff9`. +- **Monitor security updates**: Resolved `CVE-2024-6345` and `CVE-2025-47273`. +- **Collector API security updates**: Resolved `CVE-2025-66566`, `CVE-2026-10050`, `CVE-2026-40973`, `CVE-2026-40983`, `CVE-2026-40984`, `CVE-2026-42198`, `CVE-2026-42579`, `CVE-2026-42583`, `CVE-2026-44249`, `CVE-2026-45416`, `CVE-2026-45674`, `CVE-2026-47691`, `CVE-2026-50010`, `CVE-2026-54291`, `CVE-2026-54512`, `CVE-2026-54513`, `CVE-2026-59901`, and `GHSA-r7wm-3cxj-wff9`. +- **Shared utilities security updates**: Resolved `CVE-2026-21441` and `CVE-2026-33154`. +- **Operator security updates**: Resolved `CVE-2026-24051`, `CVE-2026-33186`, `CVE-2026-35469`, `CVE-2026-39821`, `CVE-2026-39829`, `CVE-2026-39830`, `CVE-2026-39831`, `CVE-2026-39832`, `CVE-2026-39833`, `CVE-2026-39834`, `CVE-2026-42508`, `CVE-2026-46595`, `CVE-2026-46597`, `CVE-2026-46680`, `CVE-2026-50151`, `CVE-2026-50163`, `CVE-2026-53488`, and `GHSA-hrxh-6v49-42gf`. +- **Collector initializer security update**: Upgraded jq to resolve `CVE-2026-32316`, `CVE-2026-33947`, `CVE-2026-33948`, `CVE-2026-39979`, `CVE-2026-40612`, `CVE-2026-41256`, `CVE-2026-41257`, `CVE-2026-43894`, `CVE-2026-43895`, `CVE-2026-43896`, `CVE-2026-44777`, `CVE-2026-47770`, `CVE-2026-49839`, and `CVE-2026-54679`. +- **Fluentd security updates**: Resolved 20 unique CVEs compared with the previous image, including 1 Critical and 11 High findings: `CVE-2025-46394`, `CVE-2025-60876`, `CVE-2025-6442`, `CVE-2026-7383`, `CVE-2026-9076`, `CVE-2026-22184`, `CVE-2026-28387`, `CVE-2026-28388`, `CVE-2026-28389`, `CVE-2026-28390`, `CVE-2026-31789`, `CVE-2026-31790`, `CVE-2026-34180`, `CVE-2026-34182`, `CVE-2026-42766`, `CVE-2026-42767`, `CVE-2026-42770`, `CVE-2026-45445`, `CVE-2026-45446`, and `CVE-2026-45447`. +- **Classic processor security updates**: Resolved `CVE-2024-6345`, `CVE-2025-47273`, and `CVE-2025-67221`. +- **API security updates**: Resolved `CVE-2024-53981`, `CVE-2024-6345`, `CVE-2025-47273`, `CVE-2025-62727`, `CVE-2026-24486`, `CVE-2026-32597`, `CVE-2026-42561`, `CVE-2026-48522`, `CVE-2026-48523`, `CVE-2026-48524`, `CVE-2026-48525`, `CVE-2026-48526`, `CVE-2026-48710`, `CVE-2026-53539`, and `CVE-2026-54283`. +- **Metrics aggregator security updates**: Resolved `CVE-2024-6345`, `CVE-2025-47273`, `CVE-2025-62727`, `CVE-2025-66418`, `CVE-2025-66471`, `CVE-2026-23490`, `CVE-2026-26007`, `CVE-2026-30922`, `CVE-2026-32597`, `CVE-2026-48710`, and `CVE-2026-54283`. From 93e57c8b1a6098c46cfbf8d43d5490fb9abe02fc Mon Sep 17 00:00:00 2001 From: Zdravko Donev Date: Wed, 16 Sep 2026 10:37:32 +0300 Subject: [PATCH 2/8] RDSC-6045: Document draining streams before Flink processor migration --- .../migration-classic-to-flink.md | 124 ++++++++++++++++-- .../installation/upgrade.md | 9 ++ 2 files changed, 120 insertions(+), 13 deletions(-) diff --git a/content/integrate/redis-data-integration/installation/migration-classic-to-flink.md b/content/integrate/redis-data-integration/installation/migration-classic-to-flink.md index 36ac9d9ba4..3f6c2a4917 100644 --- a/content/integrate/redis-data-integration/installation/migration-classic-to-flink.md +++ b/content/integrate/redis-data-integration/installation/migration-classic-to-flink.md @@ -16,7 +16,7 @@ type: integration weight: 35 --- -RDI ships with two stream processor implementations. The default *classic* +RDI ships with two stream processor implementations. The *classic* processor is implemented in Python. The *Flink* processor is built on top of [Apache Flink](https://flink.apache.org/). Both run on VM and Kubernetes installations. The Flink processor can achieve much higher throughput @@ -24,6 +24,9 @@ during snapshots, scales horizontally by changing the number of TaskManager repl and uses Flink checkpointing for fault tolerance. See [Stream processor implementations]({{< relref "/integrate/redis-data-integration/architecture#stream-processor-implementations" >}}) for an overview. +The classic processor is the default in RDI 1.19.0. The Flink processor is the +default starting with RDI 2.0.0. Select the processor explicitly when migrating. + This page describes how to migrate an existing pipeline from the classic processor to the Flink processor. The steps are the same on VMs and Kubernetes, except for the optional Helm-level tuning in [Step 1](#step-1-configure-the-flink-processor-at-the-helm-chart-level-kubernetes), @@ -31,6 +34,17 @@ which applies to Kubernetes only. ## Before you migrate +For an upgrade from RDI 1.19.0 to 2.0.0, complete the processor migration on +1.19.0 before running the 2.0.0 installer. Save your existing configuration +and jobs, and wait for the initial snapshot to finish. + +{{< warning >}} +Switching processors with records still in the RDI input streams can leave +records unprocessed. Stop collection and let the classic processor empty +the streams before switching. Zero consumer-group pending records or lag +does not prove that a stream is empty. +{{< /warning >}} + Confirm that your pipeline is compatible with the Flink processor: - `JSON.MERGE` semantics differ from the classic processor's Lua-based merge @@ -46,7 +60,7 @@ Confirm that your pipeline is compatible with the Flink processor: ## Step 1: Configure the Flink processor at the Helm chart level (Kubernetes) This step applies to **Kubernetes** installations only. On VM installations, -skip it and enable the Flink processor per pipeline in step 2. +continue with [Step 2](#step-2-disable-source-collection). The Flink processor is always available — no opt-in is required at the Helm chart level. The defaults are sized for typical workloads, so you can skip @@ -55,27 +69,110 @@ TaskManager defaults, add an `operator.dataPlane.flinkProcessor` block to your `rdi-values.yaml` file and run `helm upgrade` as described in [Configure the Flink processor]({{< relref "/integrate/redis-data-integration/installation/install-k8s#configure-the-flink-processor" >}}). Existing pipelines continue to run on the classic processor until you switch -them in step 2. +them in [Step 4](#step-4-switch-processors-and-resume-collection). For VM installations, skip this step. You can configure per-pipeline Flink -resources in step 4. +resources in [Step 7](#step-7-tune-the-flink-processor-optional). + +## Step 2: Disable source collection + +In your existing `config.yaml`, add `active: false` under the source and +keep `processors.type` set to `classic`. Use your existing source name and +preserve all other source, target, processor, and job settings. This example +shows only the fields to change: + +```yaml +sources: + : + active: false +processors: + type: classic +``` + +Deploy the complete configuration directory, including the existing jobs: + +```bash +redis-di deploy default --dir +``` + +Replace `default` if your pipeline has a different name. Wait for the +deployment to finish and the source collector to stop. Keep the pipeline +active so the classic processor can process the remaining input records. +Do not use `redis-di stop` for this step, because it also stops the processor. + +Applications can continue writing to the source database while collection +is disabled. Ensure that its change logs retain all changes for the entire +pause, so collection can resume from the saved position. + +## Step 3: Wait for the input streams to empty -## Step 2: Switch the pipeline to the Flink processor +Connect an authenticated Redis client to the **RDI database** that stores +the pipeline's input streams. Check the actual stream lengths, rather than +the target database or the consumer-group counters. -In the pipeline's `config.yaml`, set +For the default pipeline on RDI 1.19.0, find its input stream keys with +[`SCAN`]({{< relref "/commands/scan" >}}): + +```text +SCAN 0 MATCH data:{rdi}:* COUNT 1000 TYPE stream +``` + +Repeat `SCAN` with the returned cursor until it returns cursor `0`. A scan +can return an empty page before it finishes. For a different pipeline, +use its input stream prefix. Do not include dead letter queue (DLQ) streams. + +For every input stream returned, run +[`XLEN`]({{< relref "/commands/xlen" >}}): + +```text +XLEN +``` + +After the collector has stopped, require every input stream to have length +`0` in three complete checks, five seconds apart. Repeat the scan in each +check and include every stream found. Missing statistics, a connection +error, or an unexpected empty stream inventory is not proof of a drain. + +If records remain, keep the classic processor selected and resolve its +processing errors before continuing. Do not delete stream entries, reset +the pipeline, or change consumer-group positions to obtain an empty count. +Check rejected records separately: empty input streams do not prove that +every record reached the target or that DLQ history will survive a processor +change. + +## Step 4: Switch processors and resume collection + +After the drain check passes, remove the source's `active: false` setting +from the existing `config.yaml` and set [`processors.type`]({{< relref "/integrate/redis-data-integration/data-pipelines/pipeline-config#processors" >}}) to `flink`: ```yaml processors: type: flink - ... ``` -Then redeploy the pipeline. The operator stops the classic processor pods -and starts the Flink JobManager and TaskManager workloads for the pipeline. +Keep the remaining configuration and jobs, then redeploy the complete +configuration directory: + +```bash +redis-di deploy default --dir +``` + +Wait for the classic processor to terminate and the Flink JobManager and +TaskManager workloads to become healthy. Confirm that collection resumes +from the saved source position and changes committed during the pause reach +the target. Verify new inserts, updates, and deletes before upgrading RDI. + +## Step 5: Upgrade to RDI 2.0.0 + +If you are upgrading from 1.19.0 to 2.0.0, follow +[Upgrading RDI]({{< relref "/integrate/redis-data-integration/installation/upgrade" >}}) +after the processor migration succeeds. Keep `processors.type: flink` +explicitly configured. After the upgrade, verify that the pipeline is +healthy and new source changes continue to reach the target. -## Step 3: Adapt deprecated and classic-only properties +## Step 6: Adapt deprecated and classic-only properties Some `processors` properties are no-ops, classic-only, or have moved to `processors.advanced` for the Flink processor. The following table lists the @@ -97,7 +194,7 @@ and the Flink processor silently ignores classic-only top-level properties, so k both top-level properties and their `processors.advanced` equivalents lets you switch back without further edits. -## Step 4: Tune the Flink processor (optional) +## Step 7: Tune the Flink processor (optional) Fine-tune the Flink processor through the `processors.advanced` section. For example: @@ -124,7 +221,7 @@ See the [`processors.advanced` reference]({{< relref "/integrate/redis-data-integration/reference/config-yaml-reference#processors" >}}) for the full set of available properties. -## Step 5: Update observability +## Step 8: Update observability The Flink processor exposes Prometheus metrics directly from the Flink JobManager and TaskManager pods. @@ -135,6 +232,7 @@ for the `ServiceMonitor` configuration and the available metrics. ## Rolling back To revert a pipeline to the classic processor, set `processors.type` back to -`classic` (or remove the property) and redeploy the pipeline. The +`classic` and redeploy the pipeline. Do not remove the property: RDI 2.0.0 +defaults to the Flink processor. The `processors.advanced` section is silently ignored by the classic processor, so you don't need to remove it before switching back. diff --git a/content/integrate/redis-data-integration/installation/upgrade.md b/content/integrate/redis-data-integration/installation/upgrade.md index da873958c0..7eed13a485 100644 --- a/content/integrate/redis-data-integration/installation/upgrade.md +++ b/content/integrate/redis-data-integration/installation/upgrade.md @@ -16,6 +16,15 @@ type: integration weight: 30 --- +{{< warning >}} +Before upgrading a classic processor pipeline from RDI 1.19.0 to 2.0.0, +follow [Migrate from the classic processor to the Flink processor]({{< relref "/integrate/redis-data-integration/installation/migration-classic-to-flink" >}}). +Disable source collection, drain the input streams, and explicitly select +the Flink processor before running the upgrade. If you intend to keep the +classic processor, explicitly set `processors.type: classic` before upgrading; +do not rely on an omitted processor type, because the default changes in 2.0.0. +{{< /warning >}} + ## Upgrading a VM installation Follow the steps below to upgrade an existing From bb55b9e6263667097ff0fd72f623e72f1ca82804 Mon Sep 17 00:00:00 2001 From: Zdravko Donev Date: Wed, 16 Sep 2026 12:26:41 +0300 Subject: [PATCH 3/8] RDSC-6045: Separate processor migration from version upgrade --- .../migration-classic-to-flink.md | 49 ++++++++++++------- .../installation/upgrade.md | 15 +++--- 2 files changed, 38 insertions(+), 26 deletions(-) diff --git a/content/integrate/redis-data-integration/installation/migration-classic-to-flink.md b/content/integrate/redis-data-integration/installation/migration-classic-to-flink.md index 3f6c2a4917..b67c47e0f5 100644 --- a/content/integrate/redis-data-integration/installation/migration-classic-to-flink.md +++ b/content/integrate/redis-data-integration/installation/migration-classic-to-flink.md @@ -34,9 +34,11 @@ which applies to Kubernetes only. ## Before you migrate -For an upgrade from RDI 1.19.0 to 2.0.0, complete the processor migration on -1.19.0 before running the 2.0.0 installer. Save your existing configuration -and jobs, and wait for the initial snapshot to finish. +This procedure changes the pipeline processor. It does not upgrade RDI. You +can complete the migration while RDI is running version 1.19.0. + +Before you start, save your existing configuration and jobs. Wait for the +initial snapshot to finish. {{< warning >}} Switching processors with records still in the RDI input streams can leave @@ -72,7 +74,7 @@ Existing pipelines continue to run on the classic processor until you switch them in [Step 4](#step-4-switch-processors-and-resume-collection). For VM installations, skip this step. You can configure per-pipeline Flink -resources in [Step 7](#step-7-tune-the-flink-processor-optional). +resources in [Step 6](#step-6-tune-the-flink-processor-optional). ## Step 2: Disable source collection @@ -110,6 +112,8 @@ Connect an authenticated Redis client to the **RDI database** that stores the pipeline's input streams. Check the actual stream lengths, rather than the target database or the consumer-group counters. +### Check with Redis commands + For the default pipeline on RDI 1.19.0, find its input stream keys with [`SCAN`]({{< relref "/commands/scan" >}}): @@ -133,6 +137,22 @@ After the collector has stopped, require every input stream to have length check and include every stream found. Missing statistics, a connection error, or an unexpected empty stream inventory is not proof of a drain. +### Check with `redis-di` + +On RDI 1.19.0, you can also use the packaged CLI: + +```bash +redis-di describe default +redis-di list-metric-collections -p default -o json +redis-di get-metric-collection -p default -o json +``` + +Replace `default` if your pipeline has a different name. The `Pending` value +for each classic processor stream in these commands is the current stream +length. It must list every input stream and agree with the `XLEN` checks. This +value is different from the pending count of a consumer group. An `XPENDING` +count or group lag of `0` is not enough to continue. + If records remain, keep the classic processor selected and resolve its processing errors before continuing. Do not delete stream entries, reset the pipeline, or change consumer-group positions to obtain an empty count. @@ -162,17 +182,10 @@ redis-di deploy default --dir Wait for the classic processor to terminate and the Flink JobManager and TaskManager workloads to become healthy. Confirm that collection resumes from the saved source position and changes committed during the pause reach -the target. Verify new inserts, updates, and deletes before upgrading RDI. - -## Step 5: Upgrade to RDI 2.0.0 - -If you are upgrading from 1.19.0 to 2.0.0, follow -[Upgrading RDI]({{< relref "/integrate/redis-data-integration/installation/upgrade" >}}) -after the processor migration succeeds. Keep `processors.type: flink` -explicitly configured. After the upgrade, verify that the pipeline is -healthy and new source changes continue to reach the target. +the target. Verify new inserts, updates, and deletes. The processor migration +is complete after these checks pass. -## Step 6: Adapt deprecated and classic-only properties +## Step 5: Adapt deprecated and classic-only properties Some `processors` properties are no-ops, classic-only, or have moved to `processors.advanced` for the Flink processor. The following table lists the @@ -194,7 +207,7 @@ and the Flink processor silently ignores classic-only top-level properties, so k both top-level properties and their `processors.advanced` equivalents lets you switch back without further edits. -## Step 7: Tune the Flink processor (optional) +## Step 6: Tune the Flink processor (optional) Fine-tune the Flink processor through the `processors.advanced` section. For example: @@ -221,7 +234,7 @@ See the [`processors.advanced` reference]({{< relref "/integrate/redis-data-integration/reference/config-yaml-reference#processors" >}}) for the full set of available properties. -## Step 8: Update observability +## Step 7: Update observability The Flink processor exposes Prometheus metrics directly from the Flink JobManager and TaskManager pods. @@ -232,7 +245,7 @@ for the `ServiceMonitor` configuration and the available metrics. ## Rolling back To revert a pipeline to the classic processor, set `processors.type` back to -`classic` and redeploy the pipeline. Do not remove the property: RDI 2.0.0 -defaults to the Flink processor. The +`classic` and redeploy the pipeline. Keep the processor type explicit so that +the result does not depend on the default for the installed RDI version. The `processors.advanced` section is silently ignored by the classic processor, so you don't need to remove it before switching back. diff --git a/content/integrate/redis-data-integration/installation/upgrade.md b/content/integrate/redis-data-integration/installation/upgrade.md index 7eed13a485..5e0f8cd712 100644 --- a/content/integrate/redis-data-integration/installation/upgrade.md +++ b/content/integrate/redis-data-integration/installation/upgrade.md @@ -16,14 +16,13 @@ type: integration weight: 30 --- -{{< warning >}} -Before upgrading a classic processor pipeline from RDI 1.19.0 to 2.0.0, -follow [Migrate from the classic processor to the Flink processor]({{< relref "/integrate/redis-data-integration/installation/migration-classic-to-flink" >}}). -Disable source collection, drain the input streams, and explicitly select -the Flink processor before running the upgrade. If you intend to keep the -classic processor, explicitly set `processors.type: classic` before upgrading; -do not rely on an omitted processor type, because the default changes in 2.0.0. -{{< /warning >}} +{{< note >}} +Changing the pipeline processor and upgrading RDI are separate procedures. +If you want to move an existing pipeline to Flink, first follow +[Migrate from the classic processor to the Flink processor]({{< relref "/integrate/redis-data-integration/installation/migration-classic-to-flink" >}}). +You can complete that procedure on RDI 1.19.0. Before upgrading to RDI 2.0.0, +set `processors.type` explicitly to `flink` or `classic` for every pipeline. +{{< /note >}} ## Upgrading a VM installation From 8000bbac5ef1a74e85560e82484f9cbc5815cfab Mon Sep 17 00:00:00 2001 From: Zdravko Donev Date: Wed, 16 Sep 2026 12:47:10 +0300 Subject: [PATCH 4/8] Apply suggestion from @stoyanr Co-authored-by: Stoyan Rachev --- .../integrate/redis-data-integration/release-notes/rdi-2-0-0.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/content/integrate/redis-data-integration/release-notes/rdi-2-0-0.md b/content/integrate/redis-data-integration/release-notes/rdi-2-0-0.md index a4a6ce3f40..04c5552616 100644 --- a/content/integrate/redis-data-integration/release-notes/rdi-2-0-0.md +++ b/content/integrate/redis-data-integration/release-notes/rdi-2-0-0.md @@ -22,7 +22,7 @@ weight: 967 - A request now applies its changes immediately rather than waiting behind earlier in-flight requests, so the most recent request wins if several are issued in quick succession without waiting for each to complete. Polling `GET /actions/{action_id}` for an action id that a later request has superseded now returns an unknown action error instead of that action's own status. - `POST /pipelines` and `PATCH /pipelines` now also validate target connectivity similarly to their API v2 counterparts, and surface collector API failures as a `422`, `502`, `503`, or `504` error response. A request that previously succeeded despite an unreachable target now fails with a validation error. - **API v1 trace endpoint returns Not Implemented**: `POST /trace/start` now immediately returns `501 Not Implemented` instead of accepting the request and queuing a trace that would never run. -- **API v2 pipeline responses omit optional fields that have no value**: The pipeline, pipeline status, create, update, patch, start, stop, and reset responses no longer carry an optional field whose value is null; the field is left out instead. This affects `status_changed_at` and an error's `remediation`. The OpenAPI schema is unchanged, since neither field was ever required, so clients generated from it are unaffected; a client that reads either key directly has to treat it as absent rather than null. +- **API v2 pipeline responses omit optional fields that have no value**: The pipeline, pipeline status, create, update, patch, start, stop, and reset responses no longer contain optional fields with null values; such fields are left out instead. This affects `status_changed_at` and an error's `remediation`. The OpenAPI schema is unchanged, since neither field was ever required, so clients generated from it are unaffected; a client that reads either key directly has to treat it as absent rather than null. - **A source can no longer set its own `topic.prefix`**: The `topic.prefix` property is now rejected in a source's `advanced.source` section, since RDI derives the topic prefix from the source name. Remove the property from your configuration and set each job's `server_name` to the derived prefix: the source name, `rdi` for the source of an upgraded single-source pipeline handled by the Debezium collector, and the Spanner instance ID for the source of such an upgraded pipeline handled by the Flink collector. - **API v2 metric collections name data streams after their table**: The `data_streams.streams` keys of a metric collection are now the source-qualified table name (`mysql.inventory.addresses`), matching what the dead-letter queue endpoints return, instead of the Redis stream name (`{rdi}:inventory.addresses`). A client that keys off the previous form has to be updated. API v1 statistics are unchanged. - **Cassandra is no longer a source database type**: `cassandra` has been removed from the source database types that `redis-di scaffold`, the configuration template endpoint, and `rdi-admin install` offer, and a source whose connection type is `cassandra` no longer passes validation. RDI never actually supported Cassandra as source, since the Debezium Cassandra connector has to run on each Cassandra node. From e3fcf45994b65e980b35559b628d13e731398ec8 Mon Sep 17 00:00:00 2001 From: Zdravko Donev Date: Wed, 16 Sep 2026 12:48:58 +0300 Subject: [PATCH 5/8] Apply batched suggestions from code review Co-authored-by: Stoyan Rachev --- .../release-notes/rdi-2-0-0.md | 27 +++++++++---------- 1 file changed, 13 insertions(+), 14 deletions(-) diff --git a/content/integrate/redis-data-integration/release-notes/rdi-2-0-0.md b/content/integrate/redis-data-integration/release-notes/rdi-2-0-0.md index 04c5552616..6d473e2508 100644 --- a/content/integrate/redis-data-integration/release-notes/rdi-2-0-0.md +++ b/content/integrate/redis-data-integration/release-notes/rdi-2-0-0.md @@ -24,7 +24,7 @@ weight: 967 - **API v1 trace endpoint returns Not Implemented**: `POST /trace/start` now immediately returns `501 Not Implemented` instead of accepting the request and queuing a trace that would never run. - **API v2 pipeline responses omit optional fields that have no value**: The pipeline, pipeline status, create, update, patch, start, stop, and reset responses no longer contain optional fields with null values; such fields are left out instead. This affects `status_changed_at` and an error's `remediation`. The OpenAPI schema is unchanged, since neither field was ever required, so clients generated from it are unaffected; a client that reads either key directly has to treat it as absent rather than null. - **A source can no longer set its own `topic.prefix`**: The `topic.prefix` property is now rejected in a source's `advanced.source` section, since RDI derives the topic prefix from the source name. Remove the property from your configuration and set each job's `server_name` to the derived prefix: the source name, `rdi` for the source of an upgraded single-source pipeline handled by the Debezium collector, and the Spanner instance ID for the source of such an upgraded pipeline handled by the Flink collector. -- **API v2 metric collections name data streams after their table**: The `data_streams.streams` keys of a metric collection are now the source-qualified table name (`mysql.inventory.addresses`), matching what the dead-letter queue endpoints return, instead of the Redis stream name (`{rdi}:inventory.addresses`). A client that keys off the previous form has to be updated. API v1 statistics are unchanged. +- **API v2 metric collections data stream names contain the qualified table name**: The `data_streams.streams` keys of a metric collection are now the source-qualified table name (`mysql.inventory.addresses`), matching what the dead-letter queue endpoints return, instead of the Redis stream name (`{rdi}:inventory.addresses`). A client that uses the previous form has to be updated. API v1 statistics are unchanged. - **Cassandra is no longer a source database type**: `cassandra` has been removed from the source database types that `redis-di scaffold`, the configuration template endpoint, and `rdi-admin install` offer, and a source whose connection type is `cassandra` no longer passes validation. RDI never actually supported Cassandra as source, since the Debezium Cassandra connector has to run on each Cassandra node. - **The Flink processor is the default processor**: A pipeline whose `processors` section does not set `type` now deploys the Flink processor instead of the classic one. Set `processors.type` to `classic` to keep deploying the classic processor. @@ -32,34 +32,33 @@ weight: 967 - **Multiple sources in one pipeline**: A pipeline can ingest data from several source databases, of the same or different types, into one Redis target. Each source has its own name, connection settings, credentials, and collector. Use API v2 or the `redis-di` CLI to manage multi-source pipelines; API v1 supports only single-source pipelines. -- **Transformation jobs can select several tables (Flink processor)**: A job's `server_name`, `db`, `schema`, and `table` source matchers now accept a list of values in addition to a single one, and an entry prefixed with `regex:` is matched as an anchored regular expression, so one job can handle many tables. An entry without the prefix is matched literally, so existing jobs are unaffected, including table names containing regular expression characters. Two jobs whose matchers select the same table are rejected. Only the Flink processor supports this syntax, so a job using it is rejected for the classic processor. -- **Database-scoped pipeline secret keys**: Pipeline secrets in API v2 and the `redis-di` CLI now use database-independent keys (`USERNAME`, `PASSWORD`, `CACERT`, `CERT`, `KEY`, `KEY_PASSWORD`) together with a `db` parameter (`--db` in the CLI) that names the database the secret belongs to: a source name, or `target`. The previous scope-prefixed keys (`SOURCE_DB_*`, `TARGET_DB_*`) remain accepted for single-source and target secrets, used without the `db` parameter. +- **Transformation jobs can select several tables (Flink processor)**: A job's `server_name`, `db`, `schema`, and `table` source selectors now accept a list of values in addition to a single one, and an entry prefixed with `regex:` is matched as an anchored regular expression, so one job can handle many tables. An entry without the prefix is matched literally, so existing jobs are unaffected, including table names containing regular expression characters. Two jobs whose selectors select the same table are rejected. Only the Flink processor supports this syntax, so a job using it is rejected for the classic processor. +- **Database-scoped pipeline secret keys**: Pipeline secrets in API v2 and the `redis-di` CLI now use database-independent keys (`USERNAME`, `PASSWORD`, `CACERT`, `CERT`, `KEY`, `KEY_PASSWORD`) together with a `db` parameter (`--db` in the CLI) that specifies the database the secret belongs to: a source name, or `target`. The previous scope-prefixed keys (`SOURCE_DB_*`, `TARGET_DB_*`) are still accepted for single-source and target secrets, used without the `db` parameter. - **Richer validation errors from API v1 pipeline configuration endpoints**: A `422` response from `/pipelines/undeploy`, `/pipelines/sources/*`, `/pipelines/targets/*`, `/pipelines/processors/*`, and `/pipelines/secret-providers/*` may now include an `errors` array with structured per-field detail, in addition to the existing `detail` message, matching the shape already returned by `POST /pipelines` and `PATCH /pipelines`. Existing clients that only read `detail` are unaffected. -- **Pipeline components report the source they are associated with**: Each collector component in an API v2 pipeline or pipeline status response now has a `source` field naming the pipeline source it is associated with. It is empty for components that are not per-source, such as the processor and the Collector API. -- **A follower installation keeps its pipelines instead of losing them on a leadership handover**: In a high-availability RDI deployment with leader election enabled, the installation that is not currently active now keeps its Pipeline resources and reports them with a new `standby` status, instead of deleting them and recreating them on the next handover. Mutating a pipeline (create, update, patch, delete, start, stop, reset, or start/stop a source) on a standby installation now returns `503 Service Unavailable`; reads keep working. -- **Scaffolded configurations use a descriptive source name**: `redis-di scaffold` and the configuration template endpoint now name the generated source after the database type (for example `mysql`) and reference the matching source-prefixed secrets (for example `${MYSQL_DB_USERNAME}`), instead of the generic `source` name with `${SOURCE_DB_*}` references, so the config, the injected environment variables, and the secrets set with `--db ` always line up. A custom name can be passed with the new `--source-name` flag (`source_name` query parameter). The silent `rdi-admin install` writes its source secrets under the same derived names and accepts an optional `sources.default.name` key; when it reuses an existing configuration instead of scaffolding one, the secrets keep the legacy `source` name unless `sources.default.name` says otherwise. It also stores source and target certificates under their canonical file names (`ca.crt`, `client.crt`, `client.key`) in the `-db-ssl` secret, matching what the API writes, instead of the names of the files they were uploaded from. +- **Pipeline components report the source they are associated with**: Each collector component in an API v2 pipeline or pipeline status response now has a `source` field specifying the pipeline source it is associated with. It is empty for components that are not per-source, such as the processor and the Collector API. +- **A follower installation keeps its pipelines instead of deleting them on leadership loss**: In a high-availability RDI deployment with leader election enabled, the installation that is not currently active now keeps its Pipeline resources and reports them with a new `standby` status, instead of deleting them and recreating them on leadership acquisition. Mutating a pipeline (create, update, patch, delete, start, stop, reset, or start/stop a source) on a standby installation now returns `503 Service Unavailable`. +- **Scaffolded configurations use a descriptive source name**: `redis-di scaffold` and the configuration template endpoint now name the generated source after the database type (for example `mysql`). A custom name can be passed with the new `--source-name` flag (`source_name` query parameter). - **TypeScript SDK published to npmjs**: The RDI API TypeScript SDK is now publicly available on npmjs as `@rdi-ui/sdk` under the MIT license, so it no longer requires access to the internal GitHub Packages registry: `npm add @rdi-ui/sdk`. - **More collector diagnostics in API v2 metric collections**: The `GET /pipelines/{name}/metric-collections` endpoint now reports the full set of numeric Debezium collector metrics per source: event counts by operation, filtered and erroneous events, queue capacity and byte usage, the last processed transaction id, MySQL and MariaDB binlog health counters and GTID set, MongoDB primary elections, and Oracle LogMiner SCN, lag, performance, and error metrics. -- **Pipeline components report their externally-reachable endpoints**: Each component in an API v2 pipeline or pipeline status response now has an `external_endpoints` field listing the URLs at which that component is reachable from outside the cluster, discovered from its Ingress resources, in addition to the existing internal `metrics_endpoints`. -- **Source mTLS without Debezium keystore settings**: A MySQL, MariaDB, or MongoDB pipeline now presents the source client certificate without setting `database.ssl.keystore`/`mongodb.ssl.keystore` and the matching password in `advanced.source`. The collector points the connector at the keystore RDI builds from the source's certificate secrets, whether stored as pipeline secrets or supplied by a secret provider, and explicit `advanced.source` settings keep overriding the derived values. -- **Reset a single source of a multi-source pipeline**: `POST /pipelines/{name}/reset` accepts an optional `source` query parameter and `redis-di reset` an optional `--source` flag, naming one existing source that is not of type `external`. The reset then deletes only that source's data, including its change streams, offsets, schema history, dead letter queue entries, statistics, deduplication state, and record counters, and leaves every other source's data intact. The whole pipeline stops while the reset runs and starts again afterwards, as it already does for full pipeline reset. -- **Removing a source from a pipeline deletes its internal RDI data**: Removing a source with `PUT` or `PATCH /pipelines/{name}` now deletes the internal RDI data that source leaves behind, including its change streams, offsets, schema history, dead letter queue entries, statistics, deduplication state, and record counters. That data used to be kept indefinitely, and only resetting the whole pipeline removed it. The other sources keep their data, and the whole pipeline stops while the removed source's data is deleted and starts again afterwards. Records already written to the target Redis database are retained. +- **Pipeline components report their externally-reachable endpoints**: Each component in an API v2 pipeline or pipeline status response now has an `external_endpoints` field listing the URLs at which that component is reachable from outside the cluster. +- **Source mTLS without Debezium keystore settings**: A MySQL, MariaDB, or MongoDB source no longer requires setting `database.ssl.keystore`/`mongodb.ssl.keystore` and the matching password in `advanced.source` for mTLS to work correctly. +- **Removing a source from a pipeline deletes its internal RDI data**: Removing a source with `PUT` or `PATCH /pipelines/{name}` now deletes its internal RDI data, including its change streams, offsets, schema history, dead letter queue entries, statistics, deduplication state, and record counters. That data used to be kept indefinitely, and only resetting the pipeline removed it. Records already written to the target Redis database are not deleted. ### Bug Fixes - **RDI database client certificates isolated from source certificates**: The Debezium collector now presents its RDI database client certificate from a dedicated keystore instead of the shared default keystore, so the source database and RDI database client identities can no longer interfere with each other under mTLS. - **Nested processor properties shown incorrectly in `redis-di describe`**: Nested objects and arrays in a processor's advanced properties are now shown as JSON instead of a raw internal format, so the `redis-di describe` output is valid and easy to read. - **Bulk insert right after a snapshot could restart the Flink processor (Flink processor)**: A bulk insert arriving just after a snapshot completed could crash the Flink job with `Have records for a split that was not registered` while the source transitioned from snapshot to CDC processing. The job recovered automatically, but the pending change records were delayed until the restarted job claimed them. The source reader now skips record batches of splits that finished during the transition instead of failing on them. -- **Stream could stop being processed after a Flink processor restart (Flink processor)**: A narrow timing window around a checkpoint could leave a stream permanently unassigned to any reader after a restart (for example following a TaskManager failure), silently halting ingestion for that stream until the pipeline was redeployed. The processor now confirms with every reader before treating a stream as no longer needed. +- **A stream could stop being processed after a Flink processor restart (Flink processor)**: A narrow timing window around a checkpoint could leave a stream permanently unassigned to any reader after a restart (for example following a TaskManager failure), silently halting ingestion for that stream until the pipeline was redeployed. The processor now confirms with every reader before treating a stream as no longer needed. - **Unreachable source database reported as a validation error**: When RDI cannot list the source tables during pipeline validation, for example because the source database refuses the connection, the failure is now reported as a `422` validation error in the `errors` array, instead of a `422`, `502`, or `504` response carrying only a `detail` message. An unavailable collector API is still reported as a `503` response so the request can be retried. -- **No configuration template for the Redis database type**: `GET /pipelines/config/templates/ingest/redis` now returns the documented `501 Not Implemented`, instead of a `200` response carrying a configuration that named `redis` as the source connection type and left the source port empty. Redis cannot be used as a pipeline source. -- **A database flavor that does not match the database type is rejected**: `GET /pipelines/config/templates/ingest/{db_type}` now returns a `422` response when `db_flavor` names a MongoDB flavor and `db_type` is not `mongodb`, instead of a `200` response carrying a configuration that ignored the flavor. +- **No configuration template for the Redis database type**: `GET /pipelines/config/templates/ingest/redis` now returns the documented `501 Not Implemented`, instead of a `200` response with an incorrect configuration. Redis cannot be used as a pipeline source. +- **A database flavor that does not match the database type is rejected**: `GET /pipelines/config/templates/ingest/{db_type}` now returns a `422` response when `db_flavor` names a MongoDB flavor and `db_type` is not `mongodb`, instead of a `200` response with an incorrect configuration. ### Improvements - **API response body logging**: The RDI API now logs response bodies alongside request payloads. Bodies are logged at DEBUG level, while request and response metadata stays at INFO, and large bodies are truncated to keep logs manageable. - **Documented scaffolded configurations**: The configuration that `redis-di scaffold`, the configuration template endpoint, and `rdi-admin install` generate now documents every property it offers, with a description above each one, links to the pipeline configuration documentation, and the default value where one exists. It also covers Snowflake and it no longer suggests properties that the source database does not support. The target port now defaults to `6379` rather than a placeholder, so the generated configuration passes schema validation before any property is filled in. The per-database example configurations that `rdi-admin install` wrote alongside it (`config.yaml.mysql.example` and its siblings) have been removed, since the generated configuration and the documentation now cover the same ground. -- **Documented transformation job template**: The job template that `GET /pipelines/jobs/templates/ingest` returns now documents the properties of the source, transform, and output blocks, with a description above each one and links to the transformation documentation, instead of naming a handful of them in a bare skeleton. It is a valid job as returned, with only the mandatory properties left uncommented. The endpoint accepts a new optional `source_name` query parameter that names the source the job reads from, matching the parameter of the configuration template endpoint. `rdi-admin install` now writes the same template to `jobs/job.yaml`, named after the source it scaffolds, in place of the example jobs it used to copy there. +- **Documented transformation job template**: The job template that `GET /pipelines/jobs/templates/ingest` returns now documents the properties of the source, transform, and output blocks, with a description above each one and links to the transformation documentation, instead of naming a handful of them in a bare skeleton. It is a valid job as returned, with only the mandatory properties left uncommented. The endpoint accepts a new optional `source_name` query parameter that specifies the source the job reads from, matching the parameter of the configuration template endpoint. `rdi-admin install` now writes the same template to `jobs/job.yaml`, named after the source it scaffolds, in place of the example jobs it used to copy there. - **Debezium collector processes captured records on four threads**: The Debezium collector now sets `record.processing.threads` to 4 rather than letting Debezium size the pool by itself, which raises snapshot throughput considerably. Lower it through a source's `advanced.source.record.processing.threads` when fewer CPUs are available to the collector. - **Faster pipeline reset**: Resetting a pipeline is faster and no longer creates a Kubernetes Job. From 8affaaa4e3f9f19a4695c8c03c2d2ad677d2fe8f Mon Sep 17 00:00:00 2001 From: Zdravko Donev Date: Wed, 16 Sep 2026 13:08:51 +0300 Subject: [PATCH 6/8] RDSC-6038: Address RDI 2.0.0 release-note review comments --- .../release-notes/rdi-2-0-0.md | 14 +++++++------- 1 file changed, 7 insertions(+), 7 deletions(-) diff --git a/content/integrate/redis-data-integration/release-notes/rdi-2-0-0.md b/content/integrate/redis-data-integration/release-notes/rdi-2-0-0.md index 6d473e2508..c213bafcf7 100644 --- a/content/integrate/redis-data-integration/release-notes/rdi-2-0-0.md +++ b/content/integrate/redis-data-integration/release-notes/rdi-2-0-0.md @@ -17,20 +17,20 @@ weight: 967 ### Breaking Changes +- **The Flink processor is the default processor**: A pipeline whose `processors` section does not set `type` now deploys the Flink processor instead of the classic one. Set `processors.type` to `classic` to keep deploying the classic processor. - **API v1 pipeline actions apply immediately instead of queuing**: - `POST /pipelines`, `PATCH /pipelines`, `/pipelines/undeploy`, `/pipelines/sources/*`, `/pipelines/targets/*`, `/pipelines/processors/*`, `/pipelines/secret-providers/*`, `/pipelines/start`, `/pipelines/stop`, `/pipelines/reset` now update the Pipeline resource directly instead of posting to an internal task queue. Kubernetes errors are reported with their original status codes, rather than a generic 500. - A request now applies its changes immediately rather than waiting behind earlier in-flight requests, so the most recent request wins if several are issued in quick succession without waiting for each to complete. Polling `GET /actions/{action_id}` for an action id that a later request has superseded now returns an unknown action error instead of that action's own status. - `POST /pipelines` and `PATCH /pipelines` now also validate target connectivity similarly to their API v2 counterparts, and surface collector API failures as a `422`, `502`, `503`, or `504` error response. A request that previously succeeded despite an unreachable target now fails with a validation error. - **API v1 trace endpoint returns Not Implemented**: `POST /trace/start` now immediately returns `501 Not Implemented` instead of accepting the request and queuing a trace that would never run. - **API v2 pipeline responses omit optional fields that have no value**: The pipeline, pipeline status, create, update, patch, start, stop, and reset responses no longer contain optional fields with null values; such fields are left out instead. This affects `status_changed_at` and an error's `remediation`. The OpenAPI schema is unchanged, since neither field was ever required, so clients generated from it are unaffected; a client that reads either key directly has to treat it as absent rather than null. -- **A source can no longer set its own `topic.prefix`**: The `topic.prefix` property is now rejected in a source's `advanced.source` section, since RDI derives the topic prefix from the source name. Remove the property from your configuration and set each job's `server_name` to the derived prefix: the source name, `rdi` for the source of an upgraded single-source pipeline handled by the Debezium collector, and the Spanner instance ID for the source of such an upgraded pipeline handled by the Flink collector. +- **`topic.prefix` is no longer accepted in `advanced.source`**: The `topic.prefix` property is now rejected in a source's `advanced.source` section, since RDI derives the topic prefix from the source name. Remove the property from your configuration and set each job's `server_name` to the derived prefix: the source name, `rdi` for the source of an upgraded single-source pipeline handled by the Debezium collector, and the Spanner instance ID for the source of such an upgraded pipeline handled by the Flink collector. - **API v2 metric collections data stream names contain the qualified table name**: The `data_streams.streams` keys of a metric collection are now the source-qualified table name (`mysql.inventory.addresses`), matching what the dead-letter queue endpoints return, instead of the Redis stream name (`{rdi}:inventory.addresses`). A client that uses the previous form has to be updated. API v1 statistics are unchanged. - **Cassandra is no longer a source database type**: `cassandra` has been removed from the source database types that `redis-di scaffold`, the configuration template endpoint, and `rdi-admin install` offer, and a source whose connection type is `cassandra` no longer passes validation. RDI never actually supported Cassandra as source, since the Debezium Cassandra connector has to run on each Cassandra node. -- **The Flink processor is the default processor**: A pipeline whose `processors` section does not set `type` now deploys the Flink processor instead of the classic one. Set `processors.type` to `classic` to keep deploying the classic processor. ### New Features -- **Multiple sources in one pipeline**: A pipeline can ingest data from several source databases, of the same or different types, into one Redis target. Each source has its own name, connection settings, credentials, and collector. Use API v2 or the `redis-di` CLI to manage multi-source pipelines; API v1 supports only single-source pipelines. +- **Multiple sources in one pipeline**: A pipeline can ingest data from several source databases, of the same or different types, into one Redis target. Each source has its own name, connection settings, credentials, and collector. Use API v2 or the `redis-di` CLI to manage multi-source pipelines; API v1 supports only single-source pipelines. To reset one source while preserving the other sources' internal data, use the optional `source` query parameter on `POST /pipelines/{name}/reset` or the `redis-di reset --source` flag. Source reset temporarily stops and restarts the whole pipeline and does not support `external` sources. See [Multiple sources in one pipeline]({{< relref "/integrate/redis-data-integration/data-pipelines/multiple-sources" >}}). - **Transformation jobs can select several tables (Flink processor)**: A job's `server_name`, `db`, `schema`, and `table` source selectors now accept a list of values in addition to a single one, and an entry prefixed with `regex:` is matched as an anchored regular expression, so one job can handle many tables. An entry without the prefix is matched literally, so existing jobs are unaffected, including table names containing regular expression characters. Two jobs whose selectors select the same table are rejected. Only the Flink processor supports this syntax, so a job using it is rejected for the classic processor. - **Database-scoped pipeline secret keys**: Pipeline secrets in API v2 and the `redis-di` CLI now use database-independent keys (`USERNAME`, `PASSWORD`, `CACERT`, `CERT`, `KEY`, `KEY_PASSWORD`) together with a `db` parameter (`--db` in the CLI) that specifies the database the secret belongs to: a source name, or `target`. The previous scope-prefixed keys (`SOURCE_DB_*`, `TARGET_DB_*`) are still accepted for single-source and target secrets, used without the `db` parameter. @@ -47,11 +47,11 @@ weight: 967 ### Bug Fixes - **RDI database client certificates isolated from source certificates**: The Debezium collector now presents its RDI database client certificate from a dedicated keystore instead of the shared default keystore, so the source database and RDI database client identities can no longer interfere with each other under mTLS. -- **Nested processor properties shown incorrectly in `redis-di describe`**: Nested objects and arrays in a processor's advanced properties are now shown as JSON instead of a raw internal format, so the `redis-di describe` output is valid and easy to read. -- **Bulk insert right after a snapshot could restart the Flink processor (Flink processor)**: A bulk insert arriving just after a snapshot completed could crash the Flink job with `Have records for a split that was not registered` while the source transitioned from snapshot to CDC processing. The job recovered automatically, but the pending change records were delayed until the restarted job claimed them. The source reader now skips record batches of splits that finished during the transition instead of failing on them. -- **A stream could stop being processed after a Flink processor restart (Flink processor)**: A narrow timing window around a checkpoint could leave a stream permanently unassigned to any reader after a restart (for example following a TaskManager failure), silently halting ingestion for that stream until the pipeline was redeployed. The processor now confirms with every reader before treating a stream as no longer needed. +- **Nested processor properties displayed correctly in `redis-di describe`**: Nested objects and arrays in a processor's advanced properties are now shown as JSON instead of a raw internal format, so the `redis-di describe` output is valid and easy to read. +- **Bulk inserts after a snapshot no longer restart the Flink processor**: Fixed an issue where a bulk insert arriving just after a snapshot could lead to a Flink job restart with `Have records for a split that was not registered` during the transition to CDC. +- **Streams continue processing after a Flink processor restart**: A narrow timing window around a checkpoint could leave a stream permanently unassigned to any reader after a restart (for example following a TaskManager failure), silently halting ingestion for that stream until the pipeline was redeployed. The processor now confirms with every reader before treating a stream as no longer needed. - **Unreachable source database reported as a validation error**: When RDI cannot list the source tables during pipeline validation, for example because the source database refuses the connection, the failure is now reported as a `422` validation error in the `errors` array, instead of a `422`, `502`, or `504` response carrying only a `detail` message. An unavailable collector API is still reported as a `503` response so the request can be retried. -- **No configuration template for the Redis database type**: `GET /pipelines/config/templates/ingest/redis` now returns the documented `501 Not Implemented`, instead of a `200` response with an incorrect configuration. Redis cannot be used as a pipeline source. +- **Redis configuration template requests return `501 Not Implemented`**: `GET /pipelines/config/templates/ingest/redis` now returns the documented `501 Not Implemented`, instead of a `200` response with an incorrect configuration. Redis cannot be used as a pipeline source. - **A database flavor that does not match the database type is rejected**: `GET /pipelines/config/templates/ingest/{db_type}` now returns a `422` response when `db_flavor` names a MongoDB flavor and `db_type` is not `mongodb`, instead of a `200` response with an incorrect configuration. ### Improvements From 3da5154fe4a4aea28e8594500b6df447c23040d6 Mon Sep 17 00:00:00 2001 From: Zdravko Donev Date: Wed, 16 Sep 2026 13:22:10 +0300 Subject: [PATCH 7/8] RDSC-6045: Simplify Flink migration and scope upgrade guidance --- .../migration-classic-to-flink.md | 110 ++++++++---------- .../installation/upgrade.md | 30 +++-- 2 files changed, 73 insertions(+), 67 deletions(-) diff --git a/content/integrate/redis-data-integration/installation/migration-classic-to-flink.md b/content/integrate/redis-data-integration/installation/migration-classic-to-flink.md index b67c47e0f5..e992796a15 100644 --- a/content/integrate/redis-data-integration/installation/migration-classic-to-flink.md +++ b/content/integrate/redis-data-integration/installation/migration-classic-to-flink.md @@ -25,7 +25,7 @@ and uses Flink checkpointing for fault tolerance. See [Stream processor implemen for an overview. The classic processor is the default in RDI 1.19.0. The Flink processor is the -default starting with RDI 2.0.0. Select the processor explicitly when migrating. +default starting with RDI 2.0.0. This page describes how to migrate an existing pipeline from the classic processor to the Flink processor. The steps are the same on VMs and Kubernetes, @@ -34,24 +34,24 @@ which applies to Kubernetes only. ## Before you migrate -This procedure changes the pipeline processor. It does not upgrade RDI. You -can complete the migration while RDI is running version 1.19.0. +This procedure migrates the pipeline processor on RDI 1.19.0. It does not +upgrade RDI. Before you start, save your existing configuration and jobs. Wait for the -initial snapshot to finish. +initial snapshot to finish. Interrupting it causes the snapshot to restart +from the beginning. {{< warning >}} Switching processors with records still in the RDI input streams can leave -records unprocessed. Stop collection and let the classic processor empty -the streams before switching. Zero consumer-group pending records or lag -does not prove that a stream is empty. +records unprocessed. Stop the collector and let the classic processor empty +the streams before switching. {{< /warning >}} Confirm that your pipeline is compatible with the Flink processor: - `JSON.MERGE` semantics differ from the classic processor's Lua-based merge when null values are involved (see - [`use_native_json_merge`]({{< relref "/integrate/redis-data-integration/reference/config-yaml-reference#processors" >}})). + [`use_native_json_merge`]({{< relref "/integrate/redis-data-integration/reference/config-yaml-reference#processors-data-processing-configuration" >}})). The Flink processor always uses the native `JSON.MERGE` command when the target database supports it. - Ensure your Kubernetes cluster or VM has enough capacity for the Flink JobManager @@ -78,87 +78,78 @@ resources in [Step 6](#step-6-tune-the-flink-processor-optional). ## Step 2: Disable source collection -In your existing `config.yaml`, add `active: false` under the source and -keep `processors.type` set to `classic`. Use your existing source name and -preserve all other source, target, processor, and job settings. This example -shows only the fields to change: +In your existing `config.yaml`, add `active: false` under the source. Use +your existing source name and preserve all other source, target, processor, +and job settings. This example shows only the field to change: ```yaml sources: : active: false -processors: - type: classic ``` Deploy the complete configuration directory, including the existing jobs: ```bash -redis-di deploy default --dir +redis-di deploy --dir ``` -Replace `default` if your pipeline has a different name. Wait for the -deployment to finish and the source collector to stop. Keep the pipeline -active so the classic processor can process the remaining input records. +Wait for the deployment to finish and the source collector to stop. Keep +the pipeline active so the classic processor can process the remaining input records. Do not use `redis-di stop` for this step, because it also stops the processor. Applications can continue writing to the source database while collection -is disabled. Ensure that its change logs retain all changes for the entire -pause, so collection can resume from the saved position. +is disabled. When the collector restarts, it resumes from the saved source +position and processes changes made during the pause. ## Step 3: Wait for the input streams to empty -Connect an authenticated Redis client to the **RDI database** that stores -the pipeline's input streams. Check the actual stream lengths, rather than -the target database or the consumer-group counters. +After the collector has stopped, monitor and wait for every input stream +to have a length of `0`. Use either of the following methods. + +### Check with `redis-di` + +Run: + +```bash +redis-di describe +``` + +In the **Statistics** table, the **Pending** value for each classic processor +stream is its current length. Repeat the command until **Pending** is `0` +for every stream. ### Check with Redis commands -For the default pipeline on RDI 1.19.0, find its input stream keys with +Connect an authenticated Redis client to the **RDI database** that stores +the pipeline's input streams, not the target database. Find the input stream +keys with [`SCAN`]({{< relref "/commands/scan" >}}): ```text SCAN 0 MATCH data:{rdi}:* COUNT 1000 TYPE stream ``` -Repeat `SCAN` with the returned cursor until it returns cursor `0`. A scan -can return an empty page before it finishes. For a different pipeline, -use its input stream prefix. Do not include dead letter queue (DLQ) streams. - -For every input stream returned, run -[`XLEN`]({{< relref "/commands/xlen" >}}): +If the returned cursor is not `0`, pass it to the next command: ```text -XLEN +SCAN MATCH data:{rdi}:* COUNT 1000 TYPE stream ``` -After the collector has stopped, require every input stream to have length -`0` in three complete checks, five seconds apart. Repeat the scan in each -check and include every stream found. Missing statistics, a connection -error, or an unexpected empty stream inventory is not proof of a drain. +Repeat with each new cursor until the returned cursor is `0`, even if an +intermediate result contains no keys. -### Check with `redis-di` - -On RDI 1.19.0, you can also use the packaged CLI: +For every input stream returned, run +[`XLEN`]({{< relref "/commands/xlen" >}}): -```bash -redis-di describe default -redis-di list-metric-collections -p default -o json -redis-di get-metric-collection -p default -o json +```text +XLEN ``` -Replace `default` if your pipeline has a different name. The `Pending` value -for each classic processor stream in these commands is the current stream -length. It must list every input stream and agree with the `XLEN` checks. This -value is different from the pending count of a consumer group. An `XPENDING` -count or group lag of `0` is not enough to continue. +Repeat `XLEN` until every input stream has a length of `0`. -If records remain, keep the classic processor selected and resolve its -processing errors before continuing. Do not delete stream entries, reset -the pipeline, or change consumer-group positions to obtain an empty count. -Check rejected records separately: empty input streams do not prove that -every record reached the target or that DLQ history will survive a processor -change. +If records remain, keep the classic processor running and resolve its +processing errors before continuing. ## Step 4: Switch processors and resume collection @@ -172,11 +163,13 @@ processors: type: flink ``` +RDI 1.19.0 requires this setting because its default processor is `classic`. + Keep the remaining configuration and jobs, then redeploy the complete configuration directory: ```bash -redis-di deploy default --dir +redis-di deploy --dir ``` Wait for the classic processor to terminate and the Flink JobManager and @@ -231,7 +224,7 @@ processors: ``` See the -[`processors.advanced` reference]({{< relref "/integrate/redis-data-integration/reference/config-yaml-reference#processors" >}}) +[`processors.advanced` reference]({{< relref "/integrate/redis-data-integration/reference/config-yaml-reference#processorsadvanced-advanced-configuration" >}}) for the full set of available properties. ## Step 7: Update observability @@ -245,7 +238,6 @@ for the `ServiceMonitor` configuration and the available metrics. ## Rolling back To revert a pipeline to the classic processor, set `processors.type` back to -`classic` and redeploy the pipeline. Keep the processor type explicit so that -the result does not depend on the default for the installed RDI version. The -`processors.advanced` section is silently ignored by the classic processor, -so you don't need to remove it before switching back. +`classic` and redeploy the pipeline. This setting is required on RDI 2.0.0, +where the default is `flink`. The classic processor silently ignores +`processors.advanced`, so you don't need to remove it before switching back. diff --git a/content/integrate/redis-data-integration/installation/upgrade.md b/content/integrate/redis-data-integration/installation/upgrade.md index 5e0f8cd712..964881e79d 100644 --- a/content/integrate/redis-data-integration/installation/upgrade.md +++ b/content/integrate/redis-data-integration/installation/upgrade.md @@ -17,13 +17,27 @@ weight: 30 --- {{< note >}} -Changing the pipeline processor and upgrading RDI are separate procedures. -If you want to move an existing pipeline to Flink, first follow -[Migrate from the classic processor to the Flink processor]({{< relref "/integrate/redis-data-integration/installation/migration-classic-to-flink" >}}). -You can complete that procedure on RDI 1.19.0. Before upgrading to RDI 2.0.0, -set `processors.type` explicitly to `flink` or `classic` for every pipeline. +Before upgrading to RDI 2.0.0, review the +[processor default change](#upgrading-to-rdi-200). {{< /note >}} +## Upgrading to RDI 2.0.0 + +RDI 2.0.0 changes the default processor from `classic` to `flink`. This default +applies when the pipeline's `config.yaml` omits `processors.type`. + +For an existing classic pipeline, choose one of these options before upgrading: + +- To migrate to Flink, first follow + [Migrate from the classic processor to the Flink processor]({{< relref "/integrate/redis-data-integration/installation/migration-classic-to-flink" >}}) + on RDI 1.19.0. This stops the collector and drains the input streams before + switching processors. Then upgrade RDI. +- To keep the classic processor, set `processors.type: classic` in the + pipeline's `config.yaml` and deploy it before upgrading. + +If your pipeline already uses `processors.type: flink`, no processor change +is needed. Continue with the upgrade instructions for your installation. + ## Upgrading a VM installation Follow the steps below to upgrade an existing @@ -189,11 +203,11 @@ The fully supported on both VM and Kubernetes installations after upgrading to RDI 1.19.0. Once the upgrade completes, it is always available — no opt-in is required, and the defaults are sized for typical workloads. -Upgrading does not change the processor used by existing pipelines, which keep -running on the classic processor until you explicitly switch them by -setting +On RDI 1.19.0, existing classic pipelines keep using that processor until +you switch them by setting [`processors.type`]({{< relref "/integrate/redis-data-integration/data-pipelines/pipeline-config#processors" >}}) to `flink` in their `config.yaml`. +RDI 2.0.0 changes this default; see [Upgrading to RDI 2.0.0](#upgrading-to-rdi-200). On Kubernetes, to override the Flink processor defaults, add an `operator.dataPlane.flinkProcessor` block to your `rdi-values.yaml` file as From c9e8a95de925f012532cbe233e44555f83cc291f Mon Sep 17 00:00:00 2001 From: Zdravko Donev Date: Wed, 16 Sep 2026 13:32:53 +0300 Subject: [PATCH 8/8] RDSC-6045: Add verified stream drain checks --- .../migration-classic-to-flink.md | 47 ++++++++++++++----- 1 file changed, 35 insertions(+), 12 deletions(-) diff --git a/content/integrate/redis-data-integration/installation/migration-classic-to-flink.md b/content/integrate/redis-data-integration/installation/migration-classic-to-flink.md index e992796a15..419c048101 100644 --- a/content/integrate/redis-data-integration/installation/migration-classic-to-flink.md +++ b/content/integrate/redis-data-integration/installation/migration-classic-to-flink.md @@ -78,14 +78,17 @@ resources in [Step 6](#step-6-tune-the-flink-processor-optional). ## Step 2: Disable source collection -In your existing `config.yaml`, add `active: false` under the source. Use -your existing source name and preserve all other source, target, processor, -and job settings. This example shows only the field to change: +In your existing `config.yaml`, add `active: false` under the source and set +`processors.type` to `classic`. Use your existing source name and preserve +all other source, target, processor, and job settings. This example shows +only the fields to change: ```yaml sources: : active: false +processors: + type: classic ``` Deploy the complete configuration directory, including the existing jobs: @@ -99,25 +102,30 @@ the pipeline active so the classic processor can process the remaining input rec Do not use `redis-di stop` for this step, because it also stops the processor. Applications can continue writing to the source database while collection -is disabled. When the collector restarts, it resumes from the saved source -position and processes changes made during the pause. +is disabled. Make sure the database change log retains the whole paused +interval. When the collector restarts, it resumes from the saved source +position and processes those changes. ## Step 3: Wait for the input streams to empty -After the collector has stopped, monitor and wait for every input stream -to have a length of `0`. Use either of the following methods. +After the collector has stopped, wait for every input stream to have a length +of `0` in three complete checks, five seconds apart. An error, missing +statistics, or an unexpectedly empty stream list does not count as `0`. +Do not include DLQ streams. Use any of the following methods. ### Check with `redis-di` Run: ```bash -redis-di describe +redis-di describe default ``` In the **Statistics** table, the **Pending** value for each classic processor -stream is its current length. Repeat the command until **Pending** is `0` -for every stream. +stream is its current length. Confirm that every input stream is listed and +that the values agree with the Redis command checks below. This **Pending** +value is different from consumer-group pending entries. `XPENDING` or group +lag of `0` alone does not prove that a stream is empty. ### Check with Redis commands @@ -146,10 +154,25 @@ For every input stream returned, run XLEN ``` -Repeat `XLEN` until every input stream has a length of `0`. +Run a complete `SCAN` and all `XLEN` commands in each of the three checks. + +### Check with Redis Insight + +1. Connect Redis Insight to the RDI database and open **Browse**. +1. Filter by the pipeline's input stream pattern. For the default pipeline, + use `data:{rdi}:*`. Confirm that all input streams are listed. +1. Open each stream, select **Stream Data**, and use the refresh button. + Confirm that **Entries** is `0`. +1. Repeat the complete inventory and entry check three times, five seconds + apart. + +You can also open the built-in **CLI** and run the `SCAN` and `XLEN` commands +shown above. The Browser and CLI results must contain the same streams and +lengths. If records remain, keep the classic processor running and resolve its -processing errors before continuing. +processing errors before continuing. Do not delete stream entries, reset the +pipeline, or move consumer-group positions to make the count reach `0`. ## Step 4: Switch processors and resume collection