Skip to content

Commit 8a2f3e2

Browse files
committed
Address the issues
1 parent 0f1b7ce commit 8a2f3e2

File tree

12 files changed

+27
-27
lines changed

12 files changed

+27
-27
lines changed

dev/run-tests

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -173,7 +173,7 @@ CURRENT_BLOCK=$BLOCK_BUILD
173173
build/mvn $HIVE_BUILD_ARGS clean package -DskipTests
174174
else
175175
echo -e "q\n" \
176-
| build/sbt $HIVE_BUILD_ARGS package assembly/assembly streaming-kafka-assembly/assembly \
176+
| build/sbt $HIVE_BUILD_ARGS package assembly/assembly streaming-kafka-assembly/assembly \
177177
| grep -v -e "info.*Resolving" -e "warn.*Merging" -e "info.*Including"
178178
fi
179179
}

external/kafka/src/main/scala/org/apache/spark/streaming/kafka/KafkaTestUtils.scala

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -126,14 +126,14 @@ private class KafkaTestUtils extends Logging {
126126
brokerReady = true
127127
}
128128

129-
/** setup thw whole embedded servers, including Zookeeper and Kafka brokers */
130-
def setupEmbeddedServers(): Unit = {
129+
/** setup the whole embedded servers, including Zookeeper and Kafka brokers */
130+
def setup(): Unit = {
131131
setupEmbeddedZookeeper()
132132
setupEmbeddedKafkaServer()
133133
}
134134

135-
/** Tear down the whole servers, including Kafka broker and Zookeeper */
136-
def teardownEmbeddedServers(): Unit = {
135+
/** Teardown the whole servers, including Kafka broker and Zookeeper */
136+
def teardown(): Unit = {
137137
brokerReady = false
138138
zkReady = false
139139

external/kafka/src/test/java/org/apache/spark/streaming/kafka/JavaDirectKafkaStreamSuite.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -46,7 +46,7 @@ public class JavaDirectKafkaStreamSuite implements Serializable {
4646
@Before
4747
public void setUp() {
4848
kafkaTestUtils = new KafkaTestUtils();
49-
kafkaTestUtils.setupEmbeddedServers();
49+
kafkaTestUtils.setup();
5050
System.clearProperty("spark.driver.port");
5151
SparkConf sparkConf = new SparkConf()
5252
.setMaster("local[4]").setAppName(this.getClass().getSimpleName());
@@ -63,7 +63,7 @@ public void tearDown() {
6363
System.clearProperty("spark.driver.port");
6464

6565
if (kafkaTestUtils != null) {
66-
kafkaTestUtils.teardownEmbeddedServers();
66+
kafkaTestUtils.teardown();
6767
kafkaTestUtils = null;
6868
}
6969
}

external/kafka/src/test/java/org/apache/spark/streaming/kafka/JavaKafkaRDDSuite.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -42,7 +42,7 @@ public class JavaKafkaRDDSuite implements Serializable {
4242
@Before
4343
public void setUp() {
4444
kafkaTestUtils = new KafkaTestUtils();
45-
kafkaTestUtils.setupEmbeddedServers();
45+
kafkaTestUtils.setup();
4646
System.clearProperty("spark.driver.port");
4747
SparkConf sparkConf = new SparkConf()
4848
.setMaster("local[4]").setAppName(this.getClass().getSimpleName());
@@ -58,7 +58,7 @@ public void tearDown() {
5858
System.clearProperty("spark.driver.port");
5959

6060
if (kafkaTestUtils != null) {
61-
kafkaTestUtils.teardownEmbeddedServers();
61+
kafkaTestUtils.teardown();
6262
kafkaTestUtils = null;
6363
}
6464
}

external/kafka/src/test/java/org/apache/spark/streaming/kafka/JavaKafkaStreamSuite.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -47,7 +47,7 @@ public class JavaKafkaStreamSuite implements Serializable {
4747
@Before
4848
public void setUp() {
4949
kafkaTestUtils = new KafkaTestUtils();
50-
kafkaTestUtils.setupEmbeddedServers();
50+
kafkaTestUtils.setup();
5151
System.clearProperty("spark.driver.port");
5252
SparkConf sparkConf = new SparkConf()
5353
.setMaster("local[4]").setAppName(this.getClass().getSimpleName());
@@ -64,7 +64,7 @@ public void tearDown() {
6464
System.clearProperty("spark.driver.port");
6565

6666
if (kafkaTestUtils != null) {
67-
kafkaTestUtils.teardownEmbeddedServers();
67+
kafkaTestUtils.teardown();
6868
kafkaTestUtils = null;
6969
}
7070
}

external/kafka/src/test/scala/org/apache/spark/streaming/kafka/DirectKafkaStreamSuite.scala

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -54,12 +54,12 @@ class DirectKafkaStreamSuite
5454

5555
override def beforeAll {
5656
kafkaTestUtils = new KafkaTestUtils
57-
kafkaTestUtils.setupEmbeddedServers()
57+
kafkaTestUtils.setup()
5858
}
5959

6060
override def afterAll {
6161
if (kafkaTestUtils != null) {
62-
kafkaTestUtils.teardownEmbeddedServers()
62+
kafkaTestUtils.teardown()
6363
kafkaTestUtils = null
6464
}
6565
}

external/kafka/src/test/scala/org/apache/spark/streaming/kafka/KafkaClusterSuite.scala

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -31,7 +31,7 @@ class KafkaClusterSuite extends FunSuite with BeforeAndAfterAll {
3131

3232
override def beforeAll() {
3333
kafkaTestUtils = new KafkaTestUtils
34-
kafkaTestUtils.setupEmbeddedServers()
34+
kafkaTestUtils.setup()
3535

3636
kafkaTestUtils.createTopic(topic)
3737
kafkaTestUtils.sendMessages(topic, Map("a" -> 1))
@@ -40,7 +40,7 @@ class KafkaClusterSuite extends FunSuite with BeforeAndAfterAll {
4040

4141
override def afterAll() {
4242
if (kafkaTestUtils != null) {
43-
kafkaTestUtils.teardownEmbeddedServers()
43+
kafkaTestUtils.teardown()
4444
kafkaTestUtils = null
4545
}
4646
}

external/kafka/src/test/scala/org/apache/spark/streaming/kafka/KafkaRDDSuite.scala

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -37,7 +37,7 @@ class KafkaRDDSuite extends FunSuite with BeforeAndAfterAll {
3737
override def beforeAll {
3838
sc = new SparkContext(sparkConf)
3939
kafkaTestUtils = new KafkaTestUtils
40-
kafkaTestUtils.setupEmbeddedServers()
40+
kafkaTestUtils.setup()
4141
}
4242

4343
override def afterAll {
@@ -47,7 +47,7 @@ class KafkaRDDSuite extends FunSuite with BeforeAndAfterAll {
4747
}
4848

4949
if (kafkaTestUtils != null) {
50-
kafkaTestUtils.teardownEmbeddedServers()
50+
kafkaTestUtils.teardown()
5151
kafkaTestUtils = null
5252
}
5353
}

external/kafka/src/test/scala/org/apache/spark/streaming/kafka/KafkaStreamSuite.scala

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -36,7 +36,7 @@ class KafkaStreamSuite extends FunSuite with Eventually with BeforeAndAfterAll {
3636

3737
override def beforeAll(): Unit = {
3838
kafkaTestUtils = new KafkaTestUtils
39-
kafkaTestUtils.setupEmbeddedServers()
39+
kafkaTestUtils.setup()
4040
}
4141

4242
override def afterAll(): Unit = {
@@ -46,7 +46,7 @@ class KafkaStreamSuite extends FunSuite with Eventually with BeforeAndAfterAll {
4646
}
4747

4848
if (kafkaTestUtils != null) {
49-
kafkaTestUtils.teardownEmbeddedServers()
49+
kafkaTestUtils.teardown()
5050
kafkaTestUtils = null
5151
}
5252
}
@@ -84,4 +84,3 @@ class KafkaStreamSuite extends FunSuite with Eventually with BeforeAndAfterAll {
8484
}
8585
}
8686
}
87-

external/kafka/src/test/scala/org/apache/spark/streaming/kafka/ReliableKafkaStreamSuite.scala

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -53,7 +53,7 @@ class ReliableKafkaStreamSuite extends FunSuite
5353

5454
override def beforeAll() : Unit = {
5555
kafkaTestUtils = new KafkaTestUtils
56-
kafkaTestUtils.setupEmbeddedServers()
56+
kafkaTestUtils.setup()
5757

5858
groupId = s"test-consumer-${Random.nextInt(10000)}"
5959
kafkaParams = Map(
@@ -72,7 +72,7 @@ class ReliableKafkaStreamSuite extends FunSuite
7272
}
7373

7474
if (kafkaTestUtils != null) {
75-
kafkaTestUtils.teardownEmbeddedServers()
75+
kafkaTestUtils.teardown()
7676
kafkaTestUtils = null
7777
}
7878
}

0 commit comments

Comments
 (0)