diff --git a/build.gradle b/build.gradle index e60c57db20978..0253801e8015c 100644 --- a/build.gradle +++ b/build.gradle @@ -980,7 +980,7 @@ project(':streams') { project(':streams:streams-scala') { println "Building project 'streams-scala' with Scala version ${versions.scala}" apply plugin: 'scala' - archivesBaseName = "kafka-streams-scala" + archivesBaseName = "kafka-streams-scala_${versions.baseScala}" dependencies { compile project(':streams') @@ -991,7 +991,6 @@ project(':streams:streams-scala') { testCompile project(':core').sourceSets.test.output testCompile project(':streams').sourceSets.test.output testCompile project(':clients').sourceSets.test.output - testCompile libs.scalaLogging testCompile libs.junit testCompile libs.scalatest diff --git a/docs/streams/developer-guide/dsl-api.html b/docs/streams/developer-guide/dsl-api.html index ce60654d3adbb..2b25072c54991 100644 --- a/docs/streams/developer-guide/dsl-api.html +++ b/docs/streams/developer-guide/dsl-api.html @@ -3191,11 +3191,13 @@

OverviewOverview textLine.toLowerCase.split("\\W+")) .groupBy((_, word) => word) .count(Materialized.as("counts-store")) @@ -3216,7 +3218,7 @@

OverviewOverview -> to -> + // Change the stream from <user> -> <region, clicks> to <region> -> <clicks> .map((_, regionWithClicks) => regionWithClicks) // Compute the total per region by summing the individual click counts per region. diff --git a/docs/streams/index.html b/docs/streams/index.html index 8992fc5e4f3f6..72e1323388503 100644 --- a/docs/streams/index.html +++ b/docs/streams/index.html @@ -255,12 +255,13 @@

Hello Kafka Streams

import java.util.concurrent.TimeUnit import org.apache.kafka.streams.kstream.Materialized +import org.apache.kafka.streams.scala.ImplicitConversions._ +import org.apache.kafka.streams.scala._ import org.apache.kafka.streams.scala.kstream._ import org.apache.kafka.streams.{KafkaStreams, StreamsConfig} object WordCountApplication extends App { import DefaultSerdes._ - import ImplicitConversions._ val config: Properties = { val p = new Properties() @@ -269,7 +270,7 @@

Hello Kafka Streams

p } - val builder: StreamsBuilder = new StreamsBuilder() + val builder: StreamsBuilder = new StreamsBuilder val textLines: KStream[String, String] = builder.stream[String, String]("TextLinesTopic") val wordCounts: KTable[String, Long] = textLines .flatMapValues(textLine => textLine.toLowerCase.split("\\W+")) @@ -281,7 +282,7 @@

Hello Kafka Streams

streams.start() sys.ShutdownHookThread { - streams.close(10, TimeUnit.SECONDS) + streams.close(10, TimeUnit.SECONDS) } } diff --git a/streams/streams-scala/src/test/scala/org/apache/kafka/streams/scala/StreamToTableJoinScalaIntegrationTestImplicitSerdes.scala b/streams/streams-scala/src/test/scala/org/apache/kafka/streams/scala/StreamToTableJoinScalaIntegrationTestImplicitSerdes.scala index 24974c40f7283..e7014319a7247 100644 --- a/streams/streams-scala/src/test/scala/org/apache/kafka/streams/scala/StreamToTableJoinScalaIntegrationTestImplicitSerdes.scala +++ b/streams/streams-scala/src/test/scala/org/apache/kafka/streams/scala/StreamToTableJoinScalaIntegrationTestImplicitSerdes.scala @@ -34,7 +34,6 @@ import org.apache.kafka.streams._ import org.apache.kafka.streams.scala.kstream._ import ImplicitConversions._ -import com.typesafe.scalalogging.LazyLogging /** * Test suite that does an example to demonstrate stream-table joins in Kafka Streams @@ -46,7 +45,7 @@ import com.typesafe.scalalogging.LazyLogging * Hence the native Java API based version is more verbose. */ class StreamToTableJoinScalaIntegrationTestImplicitSerdes extends JUnitSuite - with StreamToTableJoinTestData with LazyLogging { + with StreamToTableJoinTestData { private val privateCluster: EmbeddedKafkaCluster = new EmbeddedKafkaCluster(1) diff --git a/streams/streams-scala/src/test/scala/org/apache/kafka/streams/scala/TopologyTest.scala b/streams/streams-scala/src/test/scala/org/apache/kafka/streams/scala/TopologyTest.scala index 89b2935b9566e..71d48341eec7f 100644 --- a/streams/streams-scala/src/test/scala/org/apache/kafka/streams/scala/TopologyTest.scala +++ b/streams/streams-scala/src/test/scala/org/apache/kafka/streams/scala/TopologyTest.scala @@ -31,7 +31,6 @@ import org.apache.kafka.streams.scala.kstream._ import org.apache.kafka.common.serialization._ import ImplicitConversions._ -import com.typesafe.scalalogging.LazyLogging import org.apache.kafka.streams.{KafkaStreams => KafkaStreamsJ, StreamsBuilder => StreamsBuilderJ, _} import org.apache.kafka.streams.kstream.{KTable => KTableJ, KStream => KStreamJ, KGroupedStream => KGroupedStreamJ, _} @@ -40,7 +39,7 @@ import collection.JavaConverters._ /** * Test suite that verifies that the topology built by the Java and Scala APIs match. */ -class TopologyTest extends JUnitSuite with LazyLogging { +class TopologyTest extends JUnitSuite { val inputTopic = "input-topic" val userClicksTopic = "user-clicks-topic" diff --git a/streams/streams-scala/src/test/scala/org/apache/kafka/streams/scala/WordCountTest.scala b/streams/streams-scala/src/test/scala/org/apache/kafka/streams/scala/WordCountTest.scala index f71e0cb4a1374..e827a3c910737 100644 --- a/streams/streams-scala/src/test/scala/org/apache/kafka/streams/scala/WordCountTest.scala +++ b/streams/streams-scala/src/test/scala/org/apache/kafka/streams/scala/WordCountTest.scala @@ -40,7 +40,6 @@ import org.apache.kafka.common.utils.MockTime import org.apache.kafka.test.TestUtils import ImplicitConversions._ -import com.typesafe.scalalogging.LazyLogging /** * Test suite that does a classic word count example. @@ -51,7 +50,7 @@ import com.typesafe.scalalogging.LazyLogging * Note: In the current project settings SAM type conversion is turned off as it's experimental in Scala 2.11. * Hence the native Java API based version is more verbose. */ -class WordCountTest extends JUnitSuite with WordCountTestData with LazyLogging { +class WordCountTest extends JUnitSuite with WordCountTestData { private val privateCluster: EmbeddedKafkaCluster = new EmbeddedKafkaCluster(1)