From 38d357bd3e118df37659c74534083ecfa39c1740 Mon Sep 17 00:00:00 2001
From: "Matthias J. Sax" Upgrading a 0.10.1
Notable changes in 0.10.2.2
+
+
+
upgrade.from added that allows rolling bounce upgrade from version 0.10.0.x Notable changes in 0.10.2.1
retries default value was changed from 0 to 10. The internal Kafka Streams consumer max.poll.interval.ms default value was changed from 300000 to Integer.MAX_VALUE.
@@ -421,6 +426,16 @@ Upgrading a 0.10.0
upgrade.from="0.10.0" set for first upgrade phase
+ (cf. KIP-268).
+
+
+ upgrade.from="0.10.0" is set for new version 0.10.1.2 upgrade.mode Notable changes in 0.10.1.0
diff --git a/gradle/dependencies.gradle b/gradle/dependencies.gradle
index cfb0b9bb39c13..daa31cd336720 100644
--- a/gradle/dependencies.gradle
+++ b/gradle/dependencies.gradle
@@ -59,6 +59,9 @@ versions += [
log4j: "1.2.17",
jopt: "5.0.4",
junit: "4.12",
+ kafka_0100: "0.10.0.1",
+ kafka_0101: "0.10.1.1",
+ kafka_0102: "0.10.2.1",
lz4: "1.4",
metrics: "2.2.0",
// PowerMock 1.x doesn't support Java 9, so use PowerMock 2.0.0 beta
@@ -95,11 +98,14 @@ libs += [
jettyServlets: "org.eclipse.jetty:jetty-servlets:$versions.jetty",
jerseyContainerServlet: "org.glassfish.jersey.containers:jersey-container-servlet:$versions.jersey",
jmhCore: "org.openjdk.jmh:jmh-core:$versions.jmh",
- jmhGeneratorAnnProcess: "org.openjdk.jmh:jmh-generator-annprocess:$versions.jmh",
jmhCoreBenchmarks: "org.openjdk.jmh:jmh-core-benchmarks:$versions.jmh",
+ jmhGeneratorAnnProcess: "org.openjdk.jmh:jmh-generator-annprocess:$versions.jmh",
+ joptSimple: "net.sf.jopt-simple:jopt-simple:$versions.jopt",
junit: "junit:junit:$versions.junit",
+ kafkaStreams_0100: "org.apache.kafka:kafka-streams:$versions.kafka_0100",
+ kafkaStreams_0101: "org.apache.kafka:kafka-streams:$versions.kafka_0101",
+ kafkaStreams_0102: "org.apache.kafka:kafka-streams:$versions.kafka_0102",
log4j: "log4j:log4j:$versions.log4j",
- joptSimple: "net.sf.jopt-simple:jopt-simple:$versions.jopt",
lz4: "org.lz4:lz4-java:$versions.lz4",
metrics: "com.yammer.metrics:metrics-core:$versions.metrics",
powermockJunit4: "org.powermock:powermock-module-junit4:$versions.powermock",
diff --git a/settings.gradle b/settings.gradle
index f0fdf07128c43..769046fe55654 100644
--- a/settings.gradle
+++ b/settings.gradle
@@ -13,5 +13,6 @@
// See the License for the specific language governing permissions and
// limitations under the License.
-include 'core', 'examples', 'clients', 'tools', 'streams', 'streams:examples', 'log4j-appender',
+include 'core', 'examples', 'clients', 'tools', 'streams', 'streams:examples', 'streams:upgrade-system-tests-0100',
+ 'streams:upgrade-system-tests-0101', 'streams:upgrade-system-tests-0102', 'log4j-appender',
'connect:api', 'connect:transforms', 'connect:runtime', 'connect:json', 'connect:file', 'jmh-benchmarks'
diff --git a/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java b/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java
index 42e65df92a6c2..ea7fc0d451d69 100644
--- a/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java
+++ b/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java
@@ -129,6 +129,11 @@ public class StreamsConfig extends AbstractConfig {
*/
public static final String PRODUCER_PREFIX = "producer.";
+ /**
+ * Config value for parameter {@link #UPGRADE_FROM_CONFIG "upgrade.from"} for upgrading an application from version {@code 0.10.0.x}.
+ */
+ public static final String UPGRADE_FROM_0100 = "0.10.0";
+
/**
* Config value for parameter {@link #PROCESSING_GUARANTEE_CONFIG "processing.guarantee"} for at-least-once processing guarantees.
*/
@@ -280,6 +285,11 @@ public class StreamsConfig extends AbstractConfig {
public static final String TIMESTAMP_EXTRACTOR_CLASS_CONFIG = "timestamp.extractor";
private static final String TIMESTAMP_EXTRACTOR_CLASS_DOC = "Timestamp extractor class that implements the org.apache.kafka.streams.processor.TimestampExtractor interface. This config is deprecated, use " + DEFAULT_TIMESTAMP_EXTRACTOR_CLASS_CONFIG + " instead";
+ /** {@code upgrade.from} */
+ public static final String UPGRADE_FROM_CONFIG = "upgrade.from";
+ public static final String UPGRADE_FROM_DOC = "Allows upgrading from version 0.10.0 to version 0.10.1 (or newer) in a backward compatible way. " +
+ "Default is null. Accepted values are \"" + UPGRADE_FROM_0100 + "\" (for upgrading from 0.10.0.x).";
+
/**
* {@code value.serde}
* @deprecated Use {@link #DEFAULT_VALUE_SERDE_CLASS_CONFIG} instead.
@@ -509,6 +519,12 @@ public class StreamsConfig extends AbstractConfig {
null,
Importance.LOW,
TIMESTAMP_EXTRACTOR_CLASS_DOC)
+ .define(UPGRADE_FROM_CONFIG,
+ ConfigDef.Type.STRING,
+ null,
+ in(null, UPGRADE_FROM_0100),
+ Importance.LOW,
+ UPGRADE_FROM_DOC)
.define(VALUE_SERDE_CLASS_CONFIG,
Type.CLASS,
null,
@@ -712,6 +728,7 @@ public Map
If you want to upgrade from 0.10.1.x to 1.0.x see the Upgrade Sections for 0.10.2, 0.11.0, and 1.0. + Note, that a brokers on-disk message format must be on version 0.10 or higher to run a Kafka Streams application version 1.0 or higher. See below a complete list of 0.10.2, 0.11.0, and 1.0 API and semantical changes that allow you to advance your application and/or simplify your code base, including the usage of new features.
Upgrading from 0.10.0.x to 1.0.x directly is also possible.
- Note, that a brokers must be on version 0.10.1 or higher to run a Kafka Streams application version 0.10.1 or higher.
+ Note, that a brokers must be on version 0.10.1 or higher and on-disk message format must be on version 0.10 or higher
+ to run a Kafka Streams application version 1.0 or higher.
See Streams API changes in 0.10.1, Streams API changes in 0.10.2,
Streams API changes in 0.11.0, and Streams API changes in 1.0
for a complete list of API changes.
diff --git a/gradle/dependencies.gradle b/gradle/dependencies.gradle
index daa31cd336720..f7027cf9de42e 100644
--- a/gradle/dependencies.gradle
+++ b/gradle/dependencies.gradle
@@ -62,6 +62,7 @@ versions += [
kafka_0100: "0.10.0.1",
kafka_0101: "0.10.1.1",
kafka_0102: "0.10.2.1",
+ kafka_0110: "0.11.0.2",
lz4: "1.4",
metrics: "2.2.0",
// PowerMock 1.x doesn't support Java 9, so use PowerMock 2.0.0 beta
@@ -105,6 +106,7 @@ libs += [
kafkaStreams_0100: "org.apache.kafka:kafka-streams:$versions.kafka_0100",
kafkaStreams_0101: "org.apache.kafka:kafka-streams:$versions.kafka_0101",
kafkaStreams_0102: "org.apache.kafka:kafka-streams:$versions.kafka_0102",
+ kafkaStreams_0110: "org.apache.kafka:kafka-streams:$versions.kafka_0110",
log4j: "log4j:log4j:$versions.log4j",
lz4: "org.lz4:lz4-java:$versions.lz4",
metrics: "com.yammer.metrics:metrics-core:$versions.metrics",
diff --git a/settings.gradle b/settings.gradle
index 769046fe55654..2820287b11edc 100644
--- a/settings.gradle
+++ b/settings.gradle
@@ -14,5 +14,6 @@
// limitations under the License.
include 'core', 'examples', 'clients', 'tools', 'streams', 'streams:examples', 'streams:upgrade-system-tests-0100',
- 'streams:upgrade-system-tests-0101', 'streams:upgrade-system-tests-0102', 'log4j-appender',
+ 'streams:upgrade-system-tests-0101', 'streams:upgrade-system-tests-0102', 'streams:upgrade-system-tests-0110',
+ 'log4j-appender',
'connect:api', 'connect:transforms', 'connect:runtime', 'connect:json', 'connect:file', 'jmh-benchmarks'
diff --git a/streams/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java b/streams/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java
index 0ee47e416ff22..5486374b62cb3 100644
--- a/streams/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java
+++ b/streams/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java
@@ -17,9 +17,9 @@
package org.apache.kafka.streams.tests;
import org.apache.kafka.streams.KafkaStreams;
+import org.apache.kafka.streams.StreamsBuilder;
import org.apache.kafka.streams.StreamsConfig;
import org.apache.kafka.streams.kstream.KStream;
-import org.apache.kafka.streams.kstream.KStreamBuilder;
import java.util.Properties;
@@ -40,8 +40,7 @@ public static void main(final String[] args) {
System.out.println("stateDir=" + stateDir);
System.out.println("upgradeFrom=" + upgradeFrom);
- final KStreamBuilder builder = new KStreamBuilder();
-
+ final StreamsBuilder builder = new StreamsBuilder();
final KStream dataStream = builder.stream("data");
dataStream.process(SmokeTestUtil.printProcessorSupplier("data"));
dataStream.to("echo");
@@ -56,7 +55,7 @@ public static void main(final String[] args) {
}
- final KafkaStreams streams = new KafkaStreams(builder, config);
+ final KafkaStreams streams = new KafkaStreams(builder.build(), config);
streams.start();
Runtime.getRuntime().addShutdownHook(new Thread() {
diff --git a/streams/upgrade-system-tests-0101/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java b/streams/upgrade-system-tests-0101/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java
index 1e6a73246d717..eebd0fab83ca0 100644
--- a/streams/upgrade-system-tests-0101/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java
+++ b/streams/upgrade-system-tests-0101/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java
@@ -30,7 +30,7 @@
public class StreamsUpgradeTest {
/**
- * This test cannot be run executed, as long as Kafka 0.10.1.2 is not release
+ * This test cannot be run executed, as long as Kafka 0.10.1.2 is not released
*/
@SuppressWarnings("unchecked")
public static void main(final String[] args) {
diff --git a/streams/upgrade-system-tests-0102/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java b/streams/upgrade-system-tests-0102/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java
index 89a12a7be5d42..18240f04ff1c5 100644
--- a/streams/upgrade-system-tests-0102/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java
+++ b/streams/upgrade-system-tests-0102/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java
@@ -30,7 +30,7 @@
public class StreamsUpgradeTest {
/**
- * This test cannot be run executed, as long as Kafka 0.10.2.2 is not release
+ * This test cannot be run executed, as long as Kafka 0.10.2.2 is not released
*/
@SuppressWarnings("unchecked")
public static void main(final String[] args) {
diff --git a/streams/upgrade-system-tests-0110/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java b/streams/upgrade-system-tests-0110/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java
new file mode 100644
index 0000000000000..779021d867245
--- /dev/null
+++ b/streams/upgrade-system-tests-0110/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java
@@ -0,0 +1,108 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.kafka.streams.tests;
+
+import org.apache.kafka.streams.KafkaStreams;
+import org.apache.kafka.streams.StreamsConfig;
+import org.apache.kafka.streams.kstream.KStream;
+import org.apache.kafka.streams.kstream.KStreamBuilder;
+import org.apache.kafka.streams.processor.AbstractProcessor;
+import org.apache.kafka.streams.processor.Processor;
+import org.apache.kafka.streams.processor.ProcessorContext;
+import org.apache.kafka.streams.processor.ProcessorSupplier;
+
+import java.util.Properties;
+
+public class StreamsUpgradeTest {
+
+ /**
+ * This test cannot be run executed, as long as Kafka 0.11.0.3 is not released
+ */
+ @SuppressWarnings("unchecked")
+ public static void main(final String[] args) {
+ if (args.length < 2) {
+ System.err.println("StreamsUpgradeTest requires three argument (kafka-url, state-dir, [upgradeFrom: optional]) but only " + args.length + " provided: "
+ + (args.length > 0 ? args[0] : ""));
+ }
+ final String kafka = args[0];
+ final String stateDir = args[1];
+ final String upgradeFrom = args.length > 2 ? args[2] : null;
+
+ System.out.println("StreamsTest instance started (StreamsUpgradeTest v0.11.0)");
+ System.out.println("kafka=" + kafka);
+ System.out.println("stateDir=" + stateDir);
+ System.out.println("upgradeFrom=" + upgradeFrom);
+
+ final KStreamBuilder builder = new KStreamBuilder();
+ final KStream dataStream = builder.stream("data");
+ dataStream.process(printProcessorSupplier());
+ dataStream.to("echo");
+
+ final Properties config = new Properties();
+ config.setProperty(StreamsConfig.APPLICATION_ID_CONFIG, "StreamsUpgradeTest");
+ config.setProperty(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, kafka);
+ config.setProperty(StreamsConfig.STATE_DIR_CONFIG, stateDir);
+ config.put(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, 1000);
+ if (upgradeFrom != null) {
+ // TODO: because Kafka 0.11.0.3 is not released yet, thus `UPGRADE_FROM_CONFIG` is not available yet
+ //config.setProperty(StreamsConfig.UPGRADE_FROM_CONFIG, upgradeFrom);
+ config.setProperty("upgrade.from", upgradeFrom);
+ }
+
+ final KafkaStreams streams = new KafkaStreams(builder, config);
+ streams.start();
+
+ Runtime.getRuntime().addShutdownHook(new Thread() {
+ @Override
+ public void run() {
+ streams.close();
+ System.out.println("UPGRADE-TEST-CLIENT-CLOSED");
+ System.out.flush();
+ }
+ });
+ }
+
+ private static Upgrade Guide & API Changes
for a complete list of API changes.
Upgrading to 1.0.2 requires two rolling bounces with config upgrade.from="0.10.0" set for first upgrade phase
(cf. KIP-268).
- As an alternative, and offline upgrade is also possible.
+ As an alternative, an offline upgrade is also possible.
upgrade.from is set to "0.10.0" for new version 1.0.2upgrade.from="0.10.0" set for first upgrade phase
(cf. KIP-268).
- As an alternative, and offline upgrade is also possible.
+ As an alternative, an offline upgrade is also possible.
upgrade.from is set to "0.10.0" for new version 0.11.0.3 upgrade.from is set to "0.10.0" for new version 0.11.0.3 upgrade.from is set to "0.10.0" for new version 0.10.2.2 upgrade.from="0.10.0" set for first upgrade phase
(cf. KIP-268).
- As an alternative, and offline upgrade is also possible.
+ As an alternative, an offline upgrade is also possible.
upgrade.from is set to "0.10.0" for new version 0.10.1.2