Versions Compared

Key

  • This line was added.
  • This line was removed.
  • Formatting was changed.

...

Code Block
languagejava
titleProduced.java
public class Produced<K, V> {
	protected Serde<K> keySerde;
	protected Serde<V> valueSerde;
	protected StreamPartitioner<? super K, ? super V> partitioner;
	protected final String topicName; // new, representing topic name
	protected final String numPartitions; // new



	public static <K, V> Produced<K, V> with(final Serde<K> keySerde,
                                         	 final Serde<V> valueSerde,
                                             final StreamPartitioner<? super K, ? super V> partitioner,
											 final topicName,	
										     final Integer numPartitions);
}

We also want to expand Grouped API with a numPartition configuration

Code Block
languagejava
titleGrouped.java
public class Grouped<K, V> {
 	protected final Serde<K> keySerde;
 	protected final Serde<V> valueSerde;
 	protected final String name;
	protected final String numPartitions; // new


	public static <K, V> Grouped<K, V> with(final String name,
    	                                    final Serde<K> keySerde,
        	                                final Serde<V> valueSerde,
											final Integer numPartitions);
}

...

 When configured with topic name and partition,  KStream application will first issue the topic lookup request and check whether the target topic is already up and running. If the target topic already exists, KStream application will log the error, and still continue the initialization.

For grouping operation,  repartition topic gets a hint on how many partitions the new topic should be created. If repartition topic is not created yet, create one with specified numPartitions, otherwise use the upstream topic partition size. If repartition logic is created, we shall issue a request to request to change the repartition topic size if current topic numPartitions != configured numPartitions. Note that there could be potential data loss, so one open question is whether we should prevent destructive operation like this.

Backward Compatibility 

Rejected Alternatives

...