On Tue, Oct 6, 2026 at 11:31 AM nitin ware <[email protected]> wrote:
> 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. > This matches some other IOs that have a "parsingFn". I'm sure there are also users that just don't bother "parsing" so of course you will want the identity function as a default, or something like that. > 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? > Nope, I think this is exactly right. And the readRows is a good idea, if it is practical. I know the Row/JSON mismatch can be severe sometimes, though typically the JSON is just as structured as a nested Row structure when you are dealing with big data sets (otherwise processing them is a nightmare). > > 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. > And possibly the interesting part here is that you might want to use a more compact encoding (even BSON if nothing else) so that you don't immediately shuffle the whole size of the JSON read. Anyhow, all of this will shake out, I'm sure. Kenn > > 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>) >>> >>
