Skip to content
This repository was archived by the owner on May 25, 2023. It is now read-only.

Commit 8f132c7

Browse files
committed
Merge branch 'develop'
2 parents 64d0688 + 8e3b56b commit 8f132c7

File tree

15 files changed

+148
-6
lines changed

15 files changed

+148
-6
lines changed

NOTICE

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,3 @@
1+
Kafka Streams Scala
2+
Copyright (C) 2018 Lightbend Inc. <https://www.lightbend.com>
3+
Copyright 2017-2018 Alexis Seigneurin.

README.md

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -16,15 +16,15 @@ The design of the library was inspired by the work started by Alexis Seigneurin
1616
`kafka-streams-scala` is published and cross-built for Scala `2.11`, and `2.12`, so you can just add the following to your build:
1717

1818
```scala
19-
val kafka_streams_scala_version = "0.2.0"
19+
val kafka_streams_scala_version = "0.2.1"
2020

2121
libraryDependencies ++= Seq("com.lightbend" %%
2222
"kafka-streams-scala" % kafka_streams_scala_version)
2323
```
2424

2525
> Note: `kafka-streams-scala` supports onwards Kafka Streams `1.0.0`.
2626
27-
The API docs for `kafka-streams-scala` is available [here](https://developer.lightbend.com/docs/api/kafka-streams-scala/0.2.0/com/lightbend/kafka/scala/streams) for Scala 2.12 and [here](https://developer.lightbend.com/docs/api/kafka-streams-scala_2.11/0.2.0/#package) for Scala 2.11.
27+
The API docs for `kafka-streams-scala` is available [here](https://developer.lightbend.com/docs/api/kafka-streams-scala/0.2.1/com/lightbend/kafka/scala/streams) for Scala 2.12 and [here](https://developer.lightbend.com/docs/api/kafka-streams-scala_2.11/0.2.1/#package) for Scala 2.11.
2828

2929
## Running the Tests
3030

build.sbt

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,7 +4,7 @@ name := "kafka-streams-scala"
44

55
organization := "com.lightbend"
66

7-
version := "0.2.0"
7+
version := "0.2.1"
88

99
scalaVersion := Versions.Scala_2_12_Version
1010

project/Versions.scala

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -7,7 +7,7 @@ object Versions {
77
val CuratorVersion = "4.0.0"
88
val MinitestVersion = "2.0.0"
99
val JDKVersion = "1.8"
10-
val Scala_2_12_Version = "2.12.4"
10+
val Scala_2_12_Version = "2.12.5"
1111
val Scala_2_11_Version = "2.11.11"
1212
val Avro4sVersion = "1.8.3"
1313
val CrossScalaVersions = Seq(Scala_2_12_Version, Scala_2_11_Version )

src/main/scala/com/lightbend/kafka/scala/streams/DefaultSerdes.scala

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
11
/**
22
* Copyright (C) 2018 Lightbend Inc. <https://www.lightbend.com>
3+
* Copyright 2017-2018 Alexis Seigneurin.
34
*/
45

56
package com.lightbend.kafka.scala.streams

src/main/scala/com/lightbend/kafka/scala/streams/FunctionConversions.scala

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
11
/**
22
* Copyright (C) 2018 Lightbend Inc. <https://www.lightbend.com>
3+
* Copyright 2017-2018 Alexis Seigneurin.
34
*/
45

56
package com.lightbend.kafka.scala.streams

src/main/scala/com/lightbend/kafka/scala/streams/ImplicitConversions.scala

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
11
/**
22
* Copyright (C) 2018 Lightbend Inc. <https://www.lightbend.com>
3+
* Copyright 2017-2018 Alexis Seigneurin.
34
*/
45

56
package com.lightbend.kafka.scala.streams

src/main/scala/com/lightbend/kafka/scala/streams/KGroupedStreamS.scala

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
11
/**
22
* Copyright (C) 2018 Lightbend Inc. <https://www.lightbend.com>
3+
* Copyright 2017-2018 Alexis Seigneurin.
34
*/
45

56
package com.lightbend.kafka.scala.streams

src/main/scala/com/lightbend/kafka/scala/streams/KGroupedTableS.scala

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
11
/**
22
* Copyright (C) 2018 Lightbend Inc. <https://www.lightbend.com>
3+
* Copyright 2017-2018 Alexis Seigneurin.
34
*/
45

56
package com.lightbend.kafka.scala.streams

src/main/scala/com/lightbend/kafka/scala/streams/KStreamS.scala

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
11
/**
22
* Copyright (C) 2018 Lightbend Inc. <https://www.lightbend.com>
3+
* Copyright 2017-2018 Alexis Seigneurin.
34
*/
45

56
package com.lightbend.kafka.scala.streams
@@ -137,7 +138,7 @@ class KStreamS[K, V](val inner: KStream[K, V]) {
137138

138139
def join[VT, VR](table: KTableS[K, VT],
139140
joiner: (V, VT) => VR)(implicit joined: Joined[K, V, VT]): KStreamS[K, VR] =
140-
inner.leftJoin[VT, VR](table.inner, joiner.asValueJoiner, joined)
141+
inner.join[VT, VR](table.inner, joiner.asValueJoiner, joined)
141142

142143
def join[GK, GV, RV](globalKTable: GlobalKTable[GK, GV],
143144
keyValueMapper: (K, V) => GK,
@@ -165,7 +166,7 @@ class KStreamS[K, V](val inner: KStream[K, V]) {
165166
windows: JoinWindows)(implicit joined: Joined[K, V, VO]): KStreamS[K, VR] =
166167
inner.outerJoin[VO, VR](otherStream.inner, joiner.asValueJoiner, windows, joined)
167168

168-
def merge(stream: KStreamS[K, V]): KStreamS[K, V] = inner.merge(stream)
169+
def merge(stream: KStreamS[K, V]): KStreamS[K, V] = inner.merge(stream.inner)
169170

170171
def peek(action: (K, V) => Unit): KStreamS[K, V] = {
171172
inner.peek(action(_,_))

0 commit comments

Comments
 (0)