Thanks Kenn for taking a look at the design doc so quickly and for the detailed feedback. This is really helpful.
> it is OK to have CouchBase SDK types in the new API, and possibly desirable. Makes sense. I'll revisit the API with that in mind. One option is for read() to accept a mapping function from the SDK's JsonObject to T, and for write() to accept functions that map T to a document ID and a JsonObject. This would also let users reuse their existing Couchbase conversion code. > The sooner you can convert a JSON string into something with a more compact encoding, the better. The mapping would happen inside the read DoFn, so the JSON representation wouldn't need to flow further through the pipeline. I'm also considering a readRows(Schema) API that converts to Row in the same DoFn. Would that be along the lines of what you had in mind, or is there another representation/pattern you'd recommend? > you may want to have .withRedistribute to optionally build in a shuffle as KafkaIO does. Agreed. I'll look at following the KafkaIO pattern with an opt-in withRedistribute() and corresponding configuration for the number of redistribution keys. I'll update the design doc based on this feedback. Thanks, Nitin Ware On Tue, 6 Oct 2026 at 09:27, Kenneth Knowles <[email protected]> wrote: > 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>) >> >
