Yes, I see. I think it should work here too. Agree to discuss details in 
Jira/PR.
Thank you very much!

> On 26 Oct 2018, at 19:13, Raghu Angadi <[email protected]> wrote:
> 
> On Fri, Oct 26, 2018 at 7:50 AM Alexey Romanenko <[email protected] 
> <mailto:[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] 
>>> <mailto:[email protected]>> wrote:
>>> 
>>> On Wed, Oct 24, 2018 at 10:47 AM Raghu Angadi <[email protected] 
>>> <mailto:[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] 
>>> <mailto:[email protected]>> wrote:
>>> I added a simple support of this for usual type of Kafka sink (PR: 
>>> https://github.com/apache/beam/pull/6776 
>>> <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] 
>>>> <mailto:[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] 
>>>> <mailto:[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] 
>>>> <mailto:[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] 
>>>> <mailto:[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 
>>>> <https://issues.apache.org/jira/browse/BEAM-5798> for the issue.
>>>> 
>>>> On Fri, Oct 19, 2018 at 2:05 PM [email protected] 
>>>> <mailto:[email protected]> <[email protected] <mailto:[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
>>>> 

Reply via email to