Thanks Kenn, that clarifies it. To summarize what I'll take into the design doc:
- read() with a parsing function, defaulting to the SDK's JsonObject when none is provided - readRows(Schema), converting to Row in the read DoFn - opt-in withRedistribute(), applied after parsing so we can avoid shuffling the raw JSON representation when a more compact encoding is available I'll share the updated doc on this thread once it's revised. Nitin On Tue, 6 Oct 2026 at 10:36, Kenneth Knowles <[email protected]> wrote: > > > 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>) >>>> >>>
