Class KafkaProcessors

java.lang.Object
com.hazelcast.jet.kafka.KafkaProcessors

public final class KafkaProcessors extends Object
Static utility class with factories of Apache Kafka source and sink processors.
Since:
Jet 3.0
  • Method Details

    • streamKafkaP

      public static <K, V, T> com.hazelcast.jet.core.ProcessorMetaSupplier streamKafkaP(@Nonnull Properties properties, @Nonnull com.hazelcast.function.FunctionEx<? super org.apache.kafka.clients.consumer.ConsumerRecord<K,V>,? extends T> projectionFn, @Nonnull com.hazelcast.jet.core.EventTimePolicy<? super T> eventTimePolicy, @Nonnull String... topics)
      Returns a supplier of processors for KafkaSources.kafka(Properties, FunctionEx, String...).
    • streamKafkaP

      public static <K, V, T> com.hazelcast.jet.core.ProcessorMetaSupplier streamKafkaP(@Nonnull Properties properties, @Nonnull com.hazelcast.function.FunctionEx<? super org.apache.kafka.clients.consumer.ConsumerRecord<K,V>,? extends T> projectionFn, @Nonnull com.hazelcast.jet.core.EventTimePolicy<? super T> eventTimePolicy, @Nonnull TopicsConfig topicsConfig)
      Returns a supplier of processors for KafkaSources.kafka(Properties, FunctionEx, TopicsConfig)}.
      Since:
      5.3
    • streamKafkaP

      public static <K, V, T> com.hazelcast.jet.core.ProcessorMetaSupplier streamKafkaP(@Nonnull com.hazelcast.jet.pipeline.DataConnectionRef dataConnectionRef, @Nonnull com.hazelcast.function.FunctionEx<? super org.apache.kafka.clients.consumer.ConsumerRecord<K,V>,? extends T> projectionFn, @Nonnull com.hazelcast.jet.core.EventTimePolicy<? super T> eventTimePolicy, @Nonnull String... topics)
    • writeKafkaP

      public static <T, K, V> com.hazelcast.jet.core.ProcessorMetaSupplier writeKafkaP(@Nonnull Properties properties, @Nonnull String topic, @Nonnull com.hazelcast.function.FunctionEx<? super T,? extends K> extractKeyFn, @Nonnull com.hazelcast.function.FunctionEx<? super T,? extends V> extractValueFn, boolean exactlyOnce)
    • writeKafkaP

      public static <T, K, V> com.hazelcast.jet.core.ProcessorMetaSupplier writeKafkaP(@Nonnull com.hazelcast.jet.pipeline.DataConnectionRef dataConnectionRef, @Nonnull String topic, @Nonnull com.hazelcast.function.FunctionEx<? super T,? extends K> extractKeyFn, @Nonnull com.hazelcast.function.FunctionEx<? super T,? extends V> extractValueFn, boolean exactlyOnce)
    • writeKafkaP

      public static <T, K, V> com.hazelcast.jet.core.ProcessorMetaSupplier writeKafkaP(@Nonnull com.hazelcast.jet.pipeline.DataConnectionRef dataConnectionRef, @Nonnull Properties properties, @Nonnull String topic, @Nonnull com.hazelcast.function.FunctionEx<? super T,? extends K> extractKeyFn, @Nonnull com.hazelcast.function.FunctionEx<? super T,? extends V> extractValueFn, boolean exactlyOnce)
    • writeKafkaP

      public static <T, K, V> com.hazelcast.jet.core.ProcessorMetaSupplier writeKafkaP(@Nonnull Properties properties, @Nonnull com.hazelcast.function.FunctionEx<? super T,? extends org.apache.kafka.clients.producer.ProducerRecord<K,V>> toRecordFn, boolean exactlyOnce)
      Returns a supplier of processors for KafkaSinks.kafka(Properties, FunctionEx).
    • writeKafkaP

      public static <T, K, V> com.hazelcast.jet.core.ProcessorMetaSupplier writeKafkaP(@Nonnull com.hazelcast.jet.pipeline.DataConnectionRef dataConnectionRef, @Nonnull com.hazelcast.function.FunctionEx<? super T,? extends org.apache.kafka.clients.producer.ProducerRecord<K,V>> toRecordFn, boolean exactlyOnce)
      Returns a supplier of processors for KafkaSinks.kafka(DataConnectionRef, FunctionEx).
    • writeKafkaP

      public static <T, K, V> com.hazelcast.jet.core.ProcessorMetaSupplier writeKafkaP(@Nonnull com.hazelcast.jet.pipeline.DataConnectionRef dataConnectionRef, @Nonnull Properties properties, @Nonnull com.hazelcast.function.FunctionEx<? super T,? extends org.apache.kafka.clients.producer.ProducerRecord<K,V>> toRecordFn, boolean exactlyOnce)
      Returns a supplier of processors for KafkaSinks.kafka(Properties, FunctionEx).