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>)
>>
>

Reply via email to