From 870167e1d33710e129506621578e019634510bc9 Mon Sep 17 00:00:00 2001 From: Jakub Scholz Date: Wed, 20 Sep 2017 13:00:45 +0200 Subject: [PATCH] Throw exception if the JSON file is invalid --- .../kafka/admin/DeleteRecordsCommand.scala | 20 +++++++++++-------- 1 file changed, 12 insertions(+), 8 deletions(-) diff --git a/core/src/main/scala/kafka/admin/DeleteRecordsCommand.scala b/core/src/main/scala/kafka/admin/DeleteRecordsCommand.scala index 2715490ec2368..a710addb9b17a 100644 --- a/core/src/main/scala/kafka/admin/DeleteRecordsCommand.scala +++ b/core/src/main/scala/kafka/admin/DeleteRecordsCommand.scala @@ -38,16 +38,20 @@ object DeleteRecordsCommand { } def parseOffsetJsonStringWithoutDedup(jsonData: String): Seq[(TopicPartition, Long)] = { - Json.parseFull(jsonData).toSeq.flatMap { js => - js.asJsonObject.get("partitions").toSeq.flatMap { partitionsJs => - partitionsJs.asJsonArray.iterator.map(_.asJsonObject).map { partitionJs => - val topic = partitionJs("topic").to[String] - val partition = partitionJs("partition").to[Int] - val offset = partitionJs("offset").to[Long] - new TopicPartition(topic, partition) -> offset - }.toBuffer + val offsetData = Json.parseFull(jsonData) + if (offsetData != None) { + offsetData.toSeq.flatMap { js => + js.asJsonObject.get("partitions").toSeq.flatMap { partitionsJs => + partitionsJs.asJsonArray.iterator.map(_.asJsonObject).map { partitionJs => + val topic = partitionJs("topic").to[String] + val partition = partitionJs("partition").to[Int] + val offset = partitionJs("offset").to[Long] + new TopicPartition(topic, partition) -> offset + }.toBuffer + } } } + else throw new AdminCommandFailedException("Offset json file doesn't contain valid JSON data.") } def execute(args: Array[String], out: PrintStream): Unit = {