Skip to content

Commit 47cce54

Browse files
scosenzajenkins
authored andcommitted
finatra-kafka-streams: Open-source Finatra Kafka Streams
Problem Finatra Kakfa Streams is an integration between Kafka Streams and Finatra which we've been using internally at Twitter for the last year. The library is not currently open-source. Solution Open-source Finatra Kafka Streams. Result Finatra Kafka Streams can be used outside of Twitter! Also, since Finatra Kafka Streams can be mixed into a Finatra HTTP or Thrift service, it is now possible to utilize Kafka Kafka Streams queryable state using a Finatra HTTP or Thrift endpoint. Differential Revision: https://phabricator.twitter.biz/D248408
1 parent 41767c6 commit 47cce54

209 files changed

Lines changed: 13689 additions & 4 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

‎CHANGELOG.rst‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -36,6 +36,10 @@ Closed
3636
Added
3737
~~~~~
3838

39+
* finatra-streams: Open-source Finatra Streams. Finatra Streams is an integration
40+
between Kafka Streams and Finatra which we've been using internally at Twitter
41+
for the last year. The library is not currently open-source.
42+
3943
Changed
4044
~~~~~~~
4145

‎build.sbt‎

Lines changed: 178 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -51,6 +51,8 @@ lazy val versions = new {
5151
// All Twitter library releases are date versioned as YY.MM.patch
5252
val twLibVersion = releaseVersion
5353

54+
val agrona = "0.9.22"
55+
val bijectionCore = "0.9.5"
5456
val commonsCodec = "1.9"
5557
val commonsFileupload = "1.3.1"
5658
val commonsIo = "2.4"
@@ -61,11 +63,13 @@ lazy val versions = new {
6163
val jodaConvert = "1.2"
6264
val jodaTime = "2.5"
6365
val junit = "4.12"
66+
val kafka = "2.0.1"
6467
val libThrift = "0.10.0"
6568
val logback = "1.1.7"
6669
val mockito = "1.9.5"
6770
val mustache = "0.8.18"
6871
val nscalaTime = "2.14.0"
72+
val rocksdbjni = "5.14.2"
6973
val scalaCheck = "1.13.4"
7074
val scalaGuice = "4.1.0"
7175
val scalaTest = "3.0.0"
@@ -213,7 +217,7 @@ lazy val finatraModules = Seq[sbt.ProjectReference](
213217
httpclient,
214218
injectApp,
215219
injectCore,
216-
injectLogback,
220+
injectLogback,
217221
injectModules,
218222
injectRequestScope,
219223
injectServer,
@@ -223,6 +227,12 @@ lazy val finatraModules = Seq[sbt.ProjectReference](
223227
injectThriftClientHttpMapper,
224228
injectUtils,
225229
jackson,
230+
kafka,
231+
kafkaStreams,
232+
kafkaStreamsPrerestore,
233+
kafkaStreamsStaticPartitioning,
234+
kafkaStreamsQueryableThriftClient,
235+
kafkaStreamsQueryableThrift,
226236
thrift,
227237
utils)
228238

@@ -754,12 +764,176 @@ lazy val injectThriftClientHttpMapper = (project in file("inject-thrift-client-h
754764
injectThriftClient % "test->test;compile->compile",
755765
thrift % "test->test;test->compile")
756766

767+
lazy val kafkaStreamsExclusionRules = Seq(
768+
ExclusionRule("javax.ws.rs", "javax.ws.rs-api"),
769+
ExclusionRule("log4j", "log4j"),
770+
ExclusionRule("org.slf4j", "slf4j-log4j12"))
771+
772+
lazy val kafkaTestJarSources =
773+
Seq("com/twitter/finatra/kafka/test/EmbeddedKafka",
774+
"com/twitter/finatra/kafka/test/KafkaTopic",
775+
"com/twitter/finatra/kafka/test/utils/ThreadUtils",
776+
"com/twitter/finatra/kafka/test/utils/PollUtils",
777+
"com/twitter/finatra/kafka/test/utils/InMemoryStatsUtil",
778+
"com/twitter/finatra/kafka/test/KafkaFeatureTest",
779+
"com/twitter/finatra/kafka/test/KafkaStateStore")
780+
lazy val kafka = (project in file("kafka"))
781+
.settings(projectSettings)
782+
.settings(
783+
name := "finatra-kafka",
784+
moduleName := "finatra-kafka",
785+
ScoverageKeys.coverageExcludedPackages := "<empty>;.*",
786+
libraryDependencies ++= Seq(
787+
"com.twitter" %% "finagle-core" % versions.twLibVersion,
788+
"com.twitter" %% "finagle-exp" % versions.twLibVersion,
789+
"com.twitter" %% "finagle-thrift" % versions.twLibVersion,
790+
"com.twitter" %% "scrooge-serializer" % versions.twLibVersion,
791+
"com.twitter" %% "util-core" % versions.twLibVersion,
792+
"org.apache.kafka" %% "kafka" % versions.kafka % "compile->compile;test->test",
793+
"org.apache.kafka" %% "kafka" % versions.kafka % "test" classifier "test",
794+
"org.apache.kafka" % "kafka-clients" % versions.kafka % "test->test",
795+
"org.apache.kafka" % "kafka-clients" % versions.kafka % "test" classifier "test",
796+
"org.apache.kafka" % "kafka-streams" % versions.kafka % "compile->compile;test->test",
797+
"org.apache.kafka" % "kafka-streams" % versions.kafka % "test" classifier "test",
798+
"org.apache.kafka" % "kafka-streams-test-utils" % versions.kafka % "compile->compile;test->test",
799+
"org.apache.kafka" % "kafka-streams-test-utils" % versions.kafka % "test" classifier "test",
800+
"org.slf4j" % "slf4j-api" % versions.slf4j % "compile->compile;test->test"
801+
),
802+
excludeDependencies in Test ++= kafkaStreamsExclusionRules,
803+
excludeDependencies ++= kafkaStreamsExclusionRules,
804+
scroogeThriftIncludeFolders in Test := Seq(file("src/test/thrift")),
805+
scroogeLanguages in Test := Seq("scala"),
806+
excludeFilter in unmanagedResources := "BUILD",
807+
publishArtifact in Test := true,
808+
mappings in (Test, packageBin) := {
809+
val previous = (mappings in (Test, packageBin)).value
810+
previous.filter(mappingContainsAnyPath(_, kafkaTestJarSources))
811+
},
812+
mappings in (Test, packageDoc) := {
813+
val previous = (mappings in (Test, packageDoc)).value
814+
previous.filter(mappingContainsAnyPath(_, kafkaTestJarSources))
815+
},
816+
mappings in (Test, packageSrc) := {
817+
val previous = (mappings in (Test, packageSrc)).value
818+
previous.filter(mappingContainsAnyPath(_, kafkaTestJarSources))
819+
}
820+
).dependsOn(
821+
injectCore % "test->test;compile->compile",
822+
injectSlf4j % "test->test;compile->compile",
823+
injectUtils % "test->test;compile->compile",
824+
jackson % "test->test",
825+
utils % "test->test;compile->compile")
826+
827+
lazy val kafkaStreamsQueryableThriftClient = (project in file("kafka-streams/kafka-streams-queryable-thrift-client"))
828+
.settings(projectSettings)
829+
.settings(
830+
name := "finatra-kafka-streams-queryable-thrift-client",
831+
moduleName := "finatra-kafka-streams-queryable-thrift-client",
832+
libraryDependencies ++= Seq(
833+
"com.twitter" %% "finagle-serversets" % versions.twLibVersion
834+
),
835+
excludeDependencies in Test ++= kafkaStreamsExclusionRules,
836+
excludeDependencies ++= kafkaStreamsExclusionRules,
837+
excludeFilter in unmanagedResources := "BUILD"
838+
).dependsOn(
839+
injectCore % "test->test;compile->compile",
840+
injectSlf4j % "test->test;compile->compile",
841+
injectUtils % "test->test;compile->compile",
842+
thrift % "test->test;compile->compile",
843+
utils % "test->test;compile->compile")
844+
845+
lazy val kafkaStreamsStaticPartitioning = (project in file("kafka-streams/kafka-streams-static-partitioning"))
846+
.settings(projectSettings)
847+
.settings(
848+
name := "finatra-kafka-streams-static-partitioning",
849+
moduleName := "finatra-kafka-streams-static-partitioning",
850+
excludeDependencies in Test ++= kafkaStreamsExclusionRules,
851+
excludeDependencies ++= kafkaStreamsExclusionRules,
852+
excludeFilter in unmanagedResources := "BUILD"
853+
).dependsOn(
854+
injectCore % "test->test;compile->compile",
855+
injectSlf4j % "test->test;compile->compile",
856+
injectUtils % "test->test;compile->compile",
857+
kafkaStreams % "test->test;compile->compile",
858+
kafkaStreamsQueryableThriftClient % "test->test;compile->compile",
859+
thrift % "test->test;compile->compile",
860+
utils % "test->test;compile->compile")
861+
862+
lazy val kafkaStreamsPrerestore = (project in file("kafka-streams/kafka-streams-prerestore"))
863+
.settings(projectSettings)
864+
.settings(
865+
name := "finatra-kafka-streams-prerestore",
866+
moduleName := "finatra-kafka-streams-prerestore",
867+
excludeDependencies in Test ++= kafkaStreamsExclusionRules,
868+
excludeDependencies ++= kafkaStreamsExclusionRules,
869+
excludeFilter in unmanagedResources := "BUILD"
870+
).dependsOn(
871+
injectCore % "test->test;compile->compile",
872+
injectSlf4j % "test->test;compile->compile",
873+
injectUtils % "test->test;compile->compile",
874+
kafkaStreams % "test->test;compile->compile",
875+
kafkaStreamsStaticPartitioning % "test->test;compile->compile",
876+
thrift % "test->test;compile->compile",
877+
utils % "test->test;compile->compile")
878+
879+
lazy val kafkaStreamsQueryableThrift = (project in file("kafka-streams/kafka-streams-queryable-thrift"))
880+
.settings(projectSettings)
881+
.settings(
882+
name := "finatra-kafka-streams-queryable-thrift",
883+
moduleName := "finatra-kafka-streams-queryable-thrift",
884+
ScoverageKeys.coverageExcludedPackages := "<empty>;.*",
885+
excludeDependencies in Test ++= kafkaStreamsExclusionRules,
886+
excludeDependencies ++= kafkaStreamsExclusionRules,
887+
scroogeThriftIncludeFolders in Compile := Seq(file("src/test/thrift")),
888+
scroogeLanguages in Compile := Seq("java", "scala"),
889+
scroogeLanguages in Test := Seq("java", "scala"),
890+
excludeFilter in unmanagedResources := "BUILD"
891+
).dependsOn(
892+
injectCore % "test->test;compile->compile",
893+
injectSlf4j % "test->test;compile->compile",
894+
injectUtils % "test->test;compile->compile",
895+
kafkaStreams % "test->test;compile->compile",
896+
kafkaStreamsQueryableThriftClient % "test->test;compile->compile",
897+
kafkaStreamsStaticPartitioning % "test->test;compile->compile",
898+
thrift % "test->test;compile->compile",
899+
utils % "test->test;compile->compile")
900+
901+
lazy val kafkaStreams = (project in file("kafka-streams/kafka-streams"))
902+
.settings(projectSettings)
903+
.settings(
904+
name := "finatra-kafka-streams",
905+
moduleName := "finatra-kafka-streams",
906+
ScoverageKeys.coverageExcludedPackages := "<empty>;.*",
907+
libraryDependencies ++= Seq(
908+
"jakarta.ws.rs" % "jakarta.ws.rs-api" % "2.1.3",
909+
"org.agrona" % "agrona" % versions.agrona,
910+
"org.apache.kafka" %% "kafka-streams-scala" % versions.kafka % "compile->compile;test->test",
911+
"org.rocksdb" % "rocksdbjni" % versions.rocksdbjni % "provided;compile->compile;test->test",
912+
"org.apache.kafka" % "kafka-streams" % versions.kafka % "compile->compile;test->test",
913+
"org.apache.kafka" % "kafka-streams" % versions.kafka % "test" classifier "test",
914+
),
915+
excludeDependencies in Test ++= kafkaStreamsExclusionRules,
916+
excludeDependencies ++= kafkaStreamsExclusionRules,
917+
excludeFilter in unmanagedResources := "BUILD",
918+
publishArtifact in Test := true
919+
).dependsOn(
920+
injectCore % "test->test;compile->compile",
921+
injectLogback % "test->test",
922+
injectSlf4j % "test->test;compile->compile",
923+
injectUtils % "test->test;compile->compile",
924+
jackson % "test->test;compile->compile",
925+
kafka % "test->test;compile->compile",
926+
kafkaStreamsQueryableThriftClient % "test->test;compile->compile",
927+
thrift % "test->test",
928+
utils % "test->test;compile->compile")
929+
930+
757931
lazy val site = (project in file("doc"))
758932
.enablePlugins(SphinxPlugin)
759933
.settings(
760-
baseSettings ++ buildSettings ++ Seq(
761-
scalacOptions in doc ++= Seq("-doc-title", "Finatra", "-doc-version", version.value),
762-
includeFilter in Sphinx := ("*.html" | "*.png" | "*.svg" | "*.js" | "*.css" | "*.gif" | "*.txt")))
934+
baseSettings ++ buildSettings ++ Seq(
935+
scalacOptions in doc ++= Seq("-doc-title", "Finatra", "-doc-version", version.value),
936+
includeFilter in Sphinx := ("*.html" | "*.png" | "*.svg" | "*.js" | "*.css" | "*.gif" | "*.txt")))
763937

764938
// START EXAMPLES
765939

‎doc/src/sphinx/user-guide/index.rst‎

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -73,6 +73,13 @@ Clients
7373

7474
- :doc:`thrift/clients`
7575

76+
Kafka Streams
77+
-------------
78+
79+
- :doc:`kafka-streams/index`
80+
- :doc:`kafka-streams/examples`
81+
- :doc:`kafka-streams/testing`
82+
7683
Testing
7784
-------
7885

@@ -124,6 +131,9 @@ Testing
124131
thrift/exceptions
125132
thrift/warmup
126133
thrift/clients
134+
kafka-streams/index
135+
kafka-streams/examples
136+
kafka-streams/testing
127137
testing/index
128138
testing/embedded
129139
testing/feature_tests
Lines changed: 55 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,55 @@
1+
.. _kafka-streams_examples:
2+
3+
Examples
4+
========
5+
6+
The `integration tests <https://github.com/twitter/finatra/blob/develop/kafka-streams/kafka-streams/src/test/scala/com/twitter/unittests/integration>`__ serve as a good collection of example Finatra Kafka Streams servers.
7+
8+
Word Count Server
9+
-----------------
10+
11+
We can build a lightweight server which counts the unique words from an input topic, storing the results in RocksDB.
12+
13+
.. code:: scala
14+
15+
class WordCountRocksDbServer extends KafkaStreamsTwitterServer {
16+
17+
override val name = "wordcount"
18+
private val countStoreName = "CountsStore"
19+
20+
override protected def configureKafkaStreams(builder: StreamsBuilder): Unit = {
21+
builder.asScala
22+
.stream[Bytes, String]("TextLinesTopic")(Consumed.`with`(Serdes.Bytes, Serdes.String))
23+
.flatMapValues(_.split(' '))
24+
.groupBy((_, word) => word)(Serialized.`with`(Serdes.String, Serdes.String))
25+
.count()(Materialized.as(countStoreName))
26+
.toStream
27+
.to("WordsWithCountsTopic")(Produced.`with`(Serdes.String, ScalaSerdes.Long))
28+
}
29+
}
30+
31+
Queryable State
32+
~~~~~~~~~~~~~~~
33+
34+
We can then expose a Thrift endpoint enabling clients to directly query the state via `interactive queries <https://kafka.apache.org/21/documentation/streams/developer-guide/interactive-queries.html>`__.
35+
36+
.. code:: scala
37+
38+
class WordCountRocksDbServer extends KafkaStreamsTwitterServer with QueryableState {
39+
40+
...
41+
42+
final override def configureThrift(router: ThriftRouter): Unit = {
43+
router
44+
.add(
45+
new WordCountQueryService(
46+
queryableFinatraKeyValueStore[String, Long](
47+
storeName = countStoreName,
48+
primaryKeySerde = Serdes.String
49+
)
50+
)
51+
)
52+
}
53+
}
54+
55+
In this example, ``WordCountQueryService`` is an underlying Thrift service.
Lines changed: 56 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,56 @@
1+
.. _kafka-streams:
2+
3+
Finatra Kafka Streams
4+
=====================
5+
6+
Finatra has native integration with `Kafka Streams <https://kafka.apache.org/documentation/streams>`__ to easily build Kafka Streams applications on top of a `TwitterServer <https://github.com/twitter/twitter-server>`__.
7+
8+
Features
9+
--------
10+
11+
- Intuitive `DSL <https://github.com/twitter/finatra/blob/develop/kafka-streams/kafka-streams/src/main/scala/com/twitter/finatra/kafkastreams/internal/utils/FinatraDslV2Implicits.scala>`__ for topology creation, compatible with the `Kafka Streams DSL <https://kafka.apache.org/21/documentation/streams/developer-guide/dsl-api.html>`__
12+
- Full Kafka Streams metric integration, exposed as `TwitterServer Metrics <https://twitter.github.io/twitter-server/Features.html#metrics>`__
13+
- `RocksDB integration <#rocksdb>`__
14+
- `Queryable State <#queryable-state>`__
15+
- `Rich testing functionality <testing.html>`__
16+
17+
Basics
18+
------
19+
20+
With `KafkaStreamsTwitterServer <https://github.com/twitter/finatra/blob/develop/kafka-streams/kafka-streams/src/main/scala/com/twitter/finatra/kafkastreams/KafkaStreamsTwitterServer.scala>`__,
21+
a fully functional service can be written by simply configuring the Kafka Streams Builder via the ``configureKafkaStreams()`` lifecycle method. See the `examples <examples.html>`__ section.
22+
23+
Transformers
24+
~~~~~~~~~~~~
25+
Implement custom `transformers <https://kafka.apache.org/21/javadoc/org/apache/kafka/streams/kstream/Transformer.html>`__ using `FinatraTransformerV2 <https://github.com/twitter/finatra/blob/develop/kafka-streams/kafka-streams/src/main/scala/com/twitter/finatra/streams/transformer/FinatraTransformerV2.scala>`__.
26+
27+
Aggregations
28+
^^^^^^^^^^^^
29+
There are several included aggregating transformers, which may be used when configuring a ``StreamsBuilder``
30+
+ ``sample``
31+
+ ``sum``
32+
+ ``compositeSum``
33+
[TODO : Add as available]
34+
- Unique counting with HyperLogLog
35+
- TopK counting
36+
37+
Stores
38+
------
39+
- Stores
40+
- Timers / TimerStores
41+
42+
RocksDB
43+
~~~~~~~
44+
In addition to using `state stores <https://kafka.apache.org/21/javadoc/org/apache/kafka/streams/state/Stores.html>`__, you may also use a RocksDB-backed store. This affords all of the advantages of using `RocksDB <https://rocksdb.org/>`__, including efficient range scans.
45+
46+
Queryable State
47+
~~~~~~~~~~~~~~~
48+
Finatra Kafka Streams supports directly querying state from a store. This can be useful for creating a service that serves data aggregated within a local Topology. You can use `static partitioning <https://github.com/twitter/finatra/blob/develop/kafka-streams/kafka-streams-static-partitioning/src/main/scala/com/twitter/finatra/streams/partitioning/StaticPartitioning.scala>`__ to query an instance deterministically known to hold a key.
49+
50+
See how queryable state is used in the following `example <examples.html#queryable-state>`__.
51+
52+
Queryable Stores
53+
^^^^^^^^^^^^^^^^
54+
- `QueryableFinatraKeyValueStore <https://github.com/twitter/finatra/blob/develop/kafka-streams/kafka-streams/src/main/scala/com/twitter/finatra/streams/query/QueryableFinatraKeyValueStore.scala>`__
55+
- `QueryableFinatraWindowStore <https://github.com/twitter/finatra/blob/develop/kafka-streams/kafka-streams/src/main/scala/com/twitter/finatra/streams/query/QueryableFinatraWindowStore.scala>`__
56+
- `QueryableFinatraCompositeWindowStore <https://github.com/twitter/finatra/blob/develop/kafka-streams/kafka-streams/src/main/scala/com/twitter/finatra/streams/query/QueryableFinatraCompositeWindowStore.scala>`__
Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,6 @@
1+
.. _kafka-streams_testing:
2+
3+
Testing
4+
=======
5+
6+
Finatra Kafka Streams includes tooling that simplifies the process of writing highly testable services. See `TopologyFeatureTest <https://github.com/twitter/finatra/blob/develop/kafka-streams/kafka-streams/src/test/scala/com/twitter/finatra/streams/tests/TopologyFeatureTest.scala>`__, which includes a `FinatraTopologyTester <https://github.com/twitter/finatra/blob/develop/kafka-streams/kafka-streams/src/test/scala/com/twitter/finatra/streams/tests/FinatraTopologyTester.scala>`__ that integrates Kafka Streams' `TopologyTestDriver <https://kafka.apache.org/21/javadoc/org/apache/kafka/streams/TopologyTestDriver.html>`__ with a `KafkaStreamsTwitterServer <https://github.com/twitter/finatra/blob/develop/kafka-streams/kafka-streams/src/main/scala/com/twitter/finatra/kafkastreams/KafkaStreamsTwitterServer.scala>`__.

‎kafka-streams/PROJECT‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,7 @@
1+
owners:
2+
- messaging-group:ldap
3+
- scosenza
4+
- dbress
5+
- adams
6+
watchers:
7+
- ds-messaging@twitter.com

0 commit comments

Comments
 (0)