Nice!

Things that jump out at me immediately:

   - it is OK to have CouchBase SDK types in the new API, and possibly
   desirable. If there are libraries that work with CouchBase SDK types, then
   you want to be able to write a DoFn that invokes them in a simple manner.
   So if the documents are already a specialized Beam type you'll have to
   convert them back. For example, if your frontend app has conversions
   between Couch documents and application-specific types that you want
   invoke in your DoFns. But if you are used to working directly with JSON
   string encoding you don't need this, of course.
   - The sooner you can convert a JSON string into something with a more
   compact encoding, the better. Ideally before you shuffle. If there is a
   common or standard way to do this, you might want to bake it in.
   - This design saves parallel reads for future work. That is fine, but of
   course for large amounts of data you will likely need it. The CouchBase
   SDK client may be able to read very fast to saturate a local buffer, but
   commits to shuffle will be more expensive and need parallelism. So then if
   parallel reads are less common, you may want to have .withRedistribute
   to optionally build in a shuffle as KafkaIO does.

Kenn

On Tue, Oct 6, 2026 at 10:02 AM nitin ware <[email protected]> wrote:

> Hi Beam Community,
>
> I've picked up https://github.com/apache/beam/issues/18381 and would like
> to add a Couchbase connector to the Java SDK. Before starting the
> implementation, I'd like to get feedback from the community on the proposed
> design and initial scope.
>
> Design doc: link
> <https://docs.google.com/document/d/1ALlXaiVSa_Dn9jJxM9MsmhEhBLPW_FP8RswN8HG33iU/edit?usp=sharing>
>
> Couchbase is an established distributed document database, but Beam
> currently has no native Couchbase connector. CouchbaseIO would provide a
> reusable integration instead of requiring users to build custom DoFns
> around the Couchbase SDK.
>
> The first version is intentionally small:
>
>    -
>
>    read() runs a user-supplied SQL++ query and streams the results as a
>    bounded PCollection<CouchbaseDocument>. The query is not partitioned
>    in v1; partitioned reads are proposed as the first follow-up.
>    -
>
>    write() upserts documents with bounded concurrency. Retrying the same
>    keyed documents results in the same end state; the connector does not claim
>    exactly-once semantics.
>    -
>
>    CouchbaseDocument is a small Beam-owned type containing a document ID
>    and JSON body, keeping Couchbase SDK types out of the public API.
>    -
>
>    The connector includes connection configuration, TLS,
>    bucket/scope/collection selection, credential redaction, unit tests, and
>    integration tests against Couchbase.
>
> The design doc also describes potential follow-up work, including
> partitioned reads, additional write modes, durability, dead-letter output,
> Row/schema support, SchemaTransform/Managed I/O, and DCP-based streaming.
>
> I've listed a few open questions in the design doc, particularly around
> the v1 scope, element type, and whether SchemaTransform/Managed I/O should
> be included in the initial implementation. I'd appreciate feedback on those
> or any other aspects of the proposed design.
>
> Please reply on this thread rather than in the design doc comments so that
> the discussion remains archived on the mailing list.
>
> Once there's some agreement on the direction, I'll start sending small
> PRs.
>
> Thanks,
> Nitin Ware (@nitinware <https://github.com/nitinware>)
>

Reply via email to