Skip to content

Commit b3806ab

Browse files
author
Andrew Or
committed
Fix test
1 parent bb80bbb commit b3806ab

File tree

1 file changed

+2
-1
lines changed

1 file changed

+2
-1
lines changed

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

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -233,7 +233,8 @@ object KafkaUtils {
233233
case (tp: TopicAndPartition, Broker(host, port)) => (tp, (host, port))
234234
}.toMap
235235
}
236-
new KafkaRDD[K, V, KD, VD, R](sc, kafkaParams, offsetRanges, leaderMap, messageHandler)
236+
val cleanedHandler = sc.clean(messageHandler)
237+
new KafkaRDD[K, V, KD, VD, R](sc, kafkaParams, offsetRanges, leaderMap, cleanedHandler)
237238
}
238239

239240
/**

0 commit comments

Comments
 (0)