On Fri, Oct 26, 2018 at 7:50 AM Alexey Romanenko <[email protected]> wrote:
> Perhaps, the optimal solution would be to have a common way to transfer a >> meta data along with key/values to KafkaIO transform. Don’t you think the >> we could use KafkaRecord (used for reading currently) for this purpose? >> > > I agree, I think this is the Right Thing to do. It addresses a few other > issues too (Support for Kafka headers while writing requires this : > BEAM-4038 <https://issues.apache.org/jira/browse/BEAM-4038>). > > > Thank you for pointing to this Jira. I see that some work had been done > there but it has not been finished yet. The idea was to use > PCollection<ProducerRecord<K, V>> instead of current PCollection<KV<K, V>>, > right? > > It sounds good for me since perfectly fits the goal of this thread > (dynamic topics) but it seems that we need to change a contract of > KafkaIO.Write to use ProducerRecord<> instead of KV<>. So, it will be the > breaking changes for user API. Do you have an idea how we can make it > back-compatible? > Interface would be pretty straight forward. See 'Write<K, V>.values()' which takes PCollection<V> instead of PCollection<KV<K, V>>. I think similar technique would wor. All the existing code works without any changes and will be a one line change for users who want to write ProducerRecords. Of course, internal implementation will have more changes since we will be carrying ProducerRecords rather than KV<K, V>. I think that is fine and safe. We can discuss more details in a follow up Jira or PR. Thanks for proposing this solution. Raghu. > > > > >> > Actually I take this back. It I don't think it coupled with output >> topic and partitions. It might just work (assuming Kafka can handle >> individual transactions spanning many topics well). >> >> Do you mean that we can just take a topic name based on KV (using >> Serialisable function or other way discussed above) and use it instead of >> current spec.getTopic() ? >> > > Yes, something on those lines. In practice it might hit practical > limitations on Kafka transaction support. E.g. if a bundle has 1000 records > going to 200 distinct topics, all of those will be written in a single > transaction. Not sure how well that would work in practice, but > theoretically it should. > > Raghu. > >> >> >> On 24 Oct 2018, at 20:01, Raghu Angadi <[email protected]> wrote: >> >> On Wed, Oct 24, 2018 at 10:47 AM Raghu Angadi <[email protected]> wrote: >> >>> My bad Alexey, I will review today. I had skimmed through the patch on >>> my phone. You are right, exactly-once sink support is not required for now. >>> >> >> >> >>> It is a quite a different beast and necessarily coupled with >>> transactions on a specific topic-partitions for correctness. >>> >> Actually I take this back. It I don't think it coupled with output topic >> and partitions. It might just work (assuming Kafka can handle individual >> transactions spanning many topics well). As you mentioned, we would still >> need to plumb it through. As such we don't know if exactly-once sink is >> being used much... (I would love to hear about it if anyone is using it). >> >> >>> >>> The primary concern is with the API. The user provides a function to map >>> an output record to its topic. We have found that such an API is usually >>> problematic. E.g. what if the record does not encode enough information >>> about topic? Say we want to select topic name based on aggregation window. >>> It might be bit more code, but simpler to let the user decide topic for >>> each record _before_ writing to the sink. E.g. it could be >>> KafkaIO.Writer<KV<topic, KV<key, value>>. >>> I wanted to think a little bit more about this, but didn't get around to >>> it. I will comment on the PR today. >>> >>> thanks for the initiative and the PR. >>> Raghu. >>> On Wed, Oct 24, 2018 at 7:03 AM Alexey Romanenko < >>> [email protected]> wrote: >>> >>>> I added a simple support of this for usual type of Kafka sink (PR: >>>> https://github.com/apache/beam/pull/6776 , welcomed for review, btw :) >>>> ) >>>> >>>> In the same time, there is another, more complicated, type of sink - >>>> EOS (Exactly Once Sink). In this case the data is partitioned among fixed >>>> number of shards and it creates one ShardWriter per shard. In its >>>> order, ShardWriter depends on Kafka topic. So, seems that in case of >>>> multiple and dynamic sink topics, we need to create new ShardWriter for >>>> every new topic per shard, >>>> >>>> Is my assumption correct or I missed/misunderstood something? >>>> >>>> On 20 Oct 2018, at 01:21, Lukasz Cwik <[email protected]> wrote: >>>> >>>> Thanks Raghu, added starter and newbie labels to the issue. >>>> >>>> On Fri, Oct 19, 2018 at 4:20 PM Raghu Angadi <[email protected]> >>>> wrote: >>>> >>>>> It will be a good starter feature for someone interested in Beam & >>>>> Kafka. Writer is very simple in Beam. It is little more than a ParDo. >>>>> >>>>> On Fri, Oct 19, 2018 at 3:37 PM Dmitry Minaev <[email protected]> >>>>> wrote: >>>>> >>>>>> Lukasz, I appreciate the quick response and filing the JIRA ticket. >>>>>> Thanks for the suggestion, unfortunately, I don't have a fixed number of >>>>>> topics. Still, we'll probably use your approach for a limited number of >>>>>> topics until the functionality is added, thank you! >>>>>> >>>>>> Thanks, >>>>>> Dmitry >>>>>> >>>>>> On Fri, Oct 19, 2018 at 2:53 PM Lukasz Cwik <[email protected]> wrote: >>>>>> >>>>>>> If there are a fixed number of topics, you could partition your >>>>>>> write by structuring your pipeline as such: >>>>>>> ParDo(PartitionByTopic) ----> KafkaIO.write(topicA) >>>>>>> \---> KafkaIO.write(topicB) >>>>>>> \---> KafkaIO.write(...) >>>>>>> >>>>>>> There is no support currently for writing to Kafka dynamically based >>>>>>> upon a destination that is part of the data. >>>>>>> I filed https://issues.apache.org/jira/browse/BEAM-5798 for the >>>>>>> issue. >>>>>>> >>>>>>> On Fri, Oct 19, 2018 at 2:05 PM [email protected] <[email protected]> >>>>>>> wrote: >>>>>>> >>>>>>>> Hi guys!! >>>>>>>> >>>>>>>> I'm trying to find a way to write to a Kafka topic using >>>>>>>> KafkaIO.write() But I need to be able to get topic name dynamically >>>>>>>> based >>>>>>>> on the data received. For example, I would like to send data for one >>>>>>>> tenant >>>>>>>> to topic "data_feed_1" and for another tenant to "topic data_feed_999". >>>>>>>> I'm coming from Flink where it's possible via >>>>>>>> KeyedSerializationSchema.getTargetTopic(). >>>>>>>> Is there anything similar in KafkaIO? >>>>>>>> >>>>>>>> Thanks, >>>>>>>> Dmitry >>>>>>>> >>>>>>> -- >>>>>> >>>>>> -- >>>>>> Dmitry >>>>>> >>>>> >>>> >> >
