Skip to content
This repository was archived by the owner on May 25, 2023. It is now read-only.
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,6 @@ package com.lightbend.kafka.scala.streams
import org.apache.kafka.streams.kstream._
import org.apache.kafka.streams.{Consumed, KeyValue}
import org.apache.kafka.common.serialization.Serde

import scala.language.implicitConversions

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,6 @@ package com.lightbend.kafka.scala.streams
import org.apache.kafka.streams.kstream._
import org.apache.kafka.streams.state.KeyValueStore
import org.apache.kafka.common.utils.Bytes
import org.apache.kafka.common.serialization.Serde
import ImplicitConversions._
import FunctionConversions._

Expand All @@ -16,19 +15,10 @@ import FunctionConversions._
*/
class KGroupedStreamS[K, V](inner: KGroupedStream[K, V]) {

def count(): KTableS[K, Long] = {
val c: KTableS[K, java.lang.Long] = inner.count()
c.mapValues[Long](Long2long _)
}
def count(): KTableS[K, Long] = inner.count().asInstanceOf[KTable[K, Long]]

def count(store: String, keySerde: Option[Serde[K]] = None): KTableS[K, Long] = {
val materialized = keySerde.foldLeft(Materialized.as[K, java.lang.Long, KeyValueStore[Bytes, Array[Byte]]](store))(
(m, serde) => m.withKeySerde(serde)
)

val c: KTableS[K, java.lang.Long] = inner.count(materialized)
c.mapValues[Long](Long2long _)
}
def count(materialized: Materialized[K, Long, KeyValueStore[Bytes, Array[Byte]]]): KTableS[K, Long] =
inner.count(materialized)

def reduce(reducer: (V, V) => V): KTableS[K, V] =
inner.reduce((v1, v2) => reducer(v1, v2))
Expand All @@ -38,18 +28,6 @@ class KGroupedStreamS[K, V](inner: KGroupedStream[K, V]) {
// works perfectly with Scala 2.12 though
inner.reduce(((v1: V, v2: V) => reducer(v1, v2)).asReducer, materialized)

def reduce(reducer: (V, V) => V, storeName: String)(implicit keySerde: Serde[K],
valueSerde: Serde[V]): KTableS[K, V] =
// need this explicit asReducer for Scala 2.11 or else the SAM conversion doesn't take place
// works perfectly with Scala 2.12 though
inner.reduce(
((v1: V, v2: V) => reducer(v1, v2)).asReducer,
Materialized
.as[K, V, KeyValueStore[Bytes, Array[Byte]]](storeName)
.withKeySerde(keySerde)
.withValueSerde(valueSerde)
)

def aggregate[VR](initializer: () => VR, aggregator: (K, V, VR) => VR): KTableS[K, VR] =
inner.aggregate(initializer.asInitializer, aggregator.asAggregator)

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,14 +15,9 @@ import FunctionConversions._
*/
class KGroupedTableS[K, V](inner: KGroupedTable[K, V]) {

type ByteArrayKVStore = KeyValueStore[Bytes, Array[Byte]]
def count(): KTableS[K, Long] = inner.count().asInstanceOf[KTable[K, Long]]

def count(): KTableS[K, Long] = {
val c: KTableS[K, java.lang.Long] = inner.count()
c.mapValues[Long](Long2long(_))
}

def count(materialized: Materialized[K, Long, ByteArrayKVStore]): KTableS[K, Long] =
def count(materialized: Materialized[K, Long, KeyValueStore[Bytes, Array[Byte]]]): KTableS[K, Long] =
inner.count(materialized)

def reduce(adder: (V, V) => V, subTractor: (V, V) => V): KTableS[K, V] =
Expand All @@ -32,7 +27,7 @@ class KGroupedTableS[K, V](inner: KGroupedTable[K, V]) {

def reduce(adder: (V, V) => V,
subtractor: (V, V) => V,
materialized: Materialized[K, V, ByteArrayKVStore]): KTableS[K, V] =
materialized: Materialized[K, V, KeyValueStore[Bytes, Array[Byte]]]): KTableS[K, V] =
// need this explicit asReducer for Scala 2.11 or else the SAM conversion doesn't take place
// works perfectly with Scala 2.12 though
inner.reduce(((v1, v2) => adder(v1, v2)).asReducer, ((v1, v2) => subtractor(v1, v2)).asReducer, materialized)
Expand All @@ -43,6 +38,6 @@ class KGroupedTableS[K, V](inner: KGroupedTable[K, V]) {
def aggregate[VR](initializer: () => VR,
adder: (K, V, VR) => VR,
subtractor: (K, V, VR) => VR,
materialized: Materialized[K, VR, ByteArrayKVStore]): KTableS[K, VR] =
materialized: Materialized[K, VR, KeyValueStore[Bytes, Array[Byte]]]): KTableS[K, VR] =
inner.aggregate(initializer.asInitializer, adder.asAggregator, subtractor.asAggregator, materialized)
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
package com.lightbend.kafka.scala.streams

import org.apache.kafka.common.serialization.Serde
import org.apache.kafka.common.utils.Bytes
import org.apache.kafka.streams.kstream.{Materialized => JavaMaterialized}
import org.apache.kafka.streams.processor.StateStore
import org.apache.kafka.streams.state._

object Materialized {

def as[K, V, S <: StateStore](storeName: String)(implicit keySerde: Serde[K],
valueSerde: Serde[V]): JavaMaterialized[K, V, S] =
JavaMaterialized.as(storeName).withKeySerde(keySerde).withValueSerde(valueSerde)

def as[K, V](
supplier: WindowBytesStoreSupplier
)(implicit keySerde: Serde[K], valueSerde: Serde[V]): JavaMaterialized[K, V, WindowStore[Bytes, Array[Byte]]] =
JavaMaterialized.as(supplier).withKeySerde(keySerde).withValueSerde(valueSerde)

def as[K, V](
supplier: SessionBytesStoreSupplier
)(implicit keySerde: Serde[K], valueSerde: Serde[V]): JavaMaterialized[K, V, SessionStore[Bytes, Array[Byte]]] =
JavaMaterialized.as(supplier).withKeySerde(keySerde).withValueSerde(valueSerde)

def as[K, V](
supplier: KeyValueBytesStoreSupplier
)(implicit keySerde: Serde[K], valueSerde: Serde[V]): JavaMaterialized[K, V, KeyValueStore[Bytes, Array[Byte]]] =
JavaMaterialized.as(supplier).withKeySerde(keySerde).withValueSerde(valueSerde)
}
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,6 @@ import org.apache.kafka.streams.kstream._
import org.apache.kafka.streams.state.SessionStore
import org.apache.kafka.common.utils.Bytes
import FunctionConversions._

import ImplicitConversions._

/**
Expand All @@ -27,10 +26,7 @@ class SessionWindowedKStreamS[K, V](val inner: SessionWindowedKStream[K, V]) {
materialized: Materialized[K, VR, SessionStore[Bytes, Array[Byte]]]): KTableS[Windowed[K], VR] =
inner.aggregate(initializer.asInitializer, aggregator.asAggregator, merger.asMerger, materialized)

def count(): KTableS[Windowed[K], Long] = {
val c: KTableS[Windowed[K], java.lang.Long] = inner.count()
c.mapValues[Long](Long2long(_))
}
def count(): KTableS[K, Long] = inner.count().asInstanceOf[KTable[K, Long]]

def count(materialized: Materialized[K, Long, SessionStore[Bytes, Array[Byte]]]): KTableS[Windowed[K], Long] =
inner.count(materialized)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@ import java.util.regex.Pattern

import com.lightbend.kafka.scala.streams.ImplicitConversions._
import org.apache.kafka.common.utils.Bytes
import org.apache.kafka.streams.kstream.{GlobalKTable, Materialized}
import org.apache.kafka.streams.kstream.{GlobalKTable, Materialized => JavaMaterialized}
import org.apache.kafka.streams.processor.{ProcessorSupplier, StateStore}
import org.apache.kafka.streams.state.{KeyValueStore, StoreBuilder}
import org.apache.kafka.streams.{Consumed, StreamsBuilder, Topology}
Expand All @@ -32,15 +32,15 @@ class StreamsBuilderS(inner: StreamsBuilder = new StreamsBuilder) {
def table[K, V](topic: String)(implicit consumed: Consumed[K, V]): KTableS[K, V] =
inner.table[K, V](topic, consumed)

def table[K, V](topic: String, materialized: Materialized[K, V, KeyValueStore[Bytes, Array[Byte]]])(
def table[K, V](topic: String, materialized: JavaMaterialized[K, V, KeyValueStore[Bytes, Array[Byte]]])(
implicit consumed: Consumed[K, V]
): KTableS[K, V] =
inner.table[K, V](topic, consumed, materialized)

def globalTable[K, V](topic: String)(implicit consumed: Consumed[K, V]): GlobalKTable[K, V] =
inner.globalTable(topic, consumed)

def globalTable[K, V](topic: String, materialized: Materialized[K, V, KeyValueStore[Bytes, Array[Byte]]])(
def globalTable[K, V](topic: String, materialized: JavaMaterialized[K, V, KeyValueStore[Bytes, Array[Byte]]])(
implicit consumed: Consumed[K, V]
): GlobalKTable[K, V] =
inner.globalTable(topic, consumed, materialized)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,6 @@ package com.lightbend.kafka.scala.streams
import org.apache.kafka.streams.kstream._
import org.apache.kafka.streams.state.WindowStore
import org.apache.kafka.common.utils.Bytes
import org.apache.kafka.common.serialization.Serde
import ImplicitConversions._
import FunctionConversions._

Expand All @@ -24,19 +23,12 @@ class TimeWindowedKStreamS[K, V](val inner: TimeWindowedKStream[K, V]) {
materialized: Materialized[K, VR, WindowStore[Bytes, Array[Byte]]]): KTableS[Windowed[K], VR] =
inner.aggregate(initializer.asInitializer, aggregator.asAggregator, materialized)

def count(): KTableS[Windowed[K], Long] = {
val c: KTableS[Windowed[K], java.lang.Long] = inner.count()
c.mapValues[Long](Long2long(_))
}

def count(store: String, keySerde: Option[Serde[K]] = None): KTableS[Windowed[K], Long] = {
val materialized = {
val m = Materialized.as[K, java.lang.Long, WindowStore[Bytes, Array[Byte]]](store)
keySerde.foldLeft(m)((m, serde) => m.withKeySerde(serde))
}
val c: KTableS[Windowed[K], java.lang.Long] = inner.count(materialized)
c.mapValues[Long](Long2long(_))
}
def count(): KTableS[Windowed[K], Long] = inner.count().asInstanceOf[KTable[Windowed[K], Long]]

def count(
materialized: Materialized[Windowed[K], Long, WindowStore[Bytes, Array[Byte]]]
): KTableS[Windowed[K], Long] =
inner.count(materialized)

def reduce(reducer: (V, V) => V): KTableS[Windowed[K], V] =
inner.reduce(reducer.asReducer)
Expand Down