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

Reply via email to