Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -57,7 +57,8 @@ class DescribeLogDirsRequestTest extends BaseRequestTest {
val log1 = servers.head.logManager.getLog(tp1).get
assertEquals(log0.size, replicaInfo0.size)
assertEquals(log1.size, replicaInfo1.size)
assertTrue(servers.head.logManager.getLog(tp0).get.logEndOffset > 0)
val logEndOffset = servers.head.logManager.getLog(tp0).get.logEndOffset
assertTrue(s"LogEndOffset '$logEndOffset' should be > 0", logEndOffset > 0)
assertEquals(servers.head.replicaManager.getLogEndOffsetLag(tp0, log0.logEndOffset, false), replicaInfo0.offsetLag)
assertEquals(servers.head.replicaManager.getLogEndOffsetLag(tp1, log1.logEndOffset, false), replicaInfo1.offsetLag)
}
Expand Down
5 changes: 1 addition & 4 deletions core/src/test/scala/unit/kafka/utils/JaasTestUtils.scala
Original file line number Diff line number Diff line change
Expand Up @@ -155,10 +155,7 @@ object JaasTestUtils {
val serviceName = "kafka"

def saslConfigs(saslProperties: Option[Properties]): Properties = {
val result = saslProperties match {
case Some(properties) => properties
case None => new Properties
}
val result = saslProperties.getOrElse(new Properties)
// IBM Kerberos module doesn't support the serviceName JAAS property, hence it needs to be
// passed as a Kafka property
if (Java.isIbmJdk && !result.contains(KafkaConfig.SaslKerberosServiceNameProp))
Expand Down
7 changes: 5 additions & 2 deletions core/src/test/scala/unit/kafka/utils/TestUtils.scala
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,7 @@ import org.apache.kafka.common.internals.Topic
import org.apache.kafka.common.network.{ListenerName, Mode}
import org.apache.kafka.common.record._
import org.apache.kafka.common.security.auth.SecurityProtocol
import org.apache.kafka.common.serialization.{ByteArrayDeserializer, ByteArraySerializer, Deserializer, Serializer}
import org.apache.kafka.common.serialization.{ByteArrayDeserializer, ByteArraySerializer, Deserializer, IntegerSerializer, Serializer}
import org.apache.kafka.common.utils.Time
import org.apache.kafka.common.utils.Utils._
import org.apache.kafka.test.{TestSslUtils, TestUtils => JTestUtils}
Expand Down Expand Up @@ -1065,7 +1065,10 @@ object TestUtils extends Logging {
numMessages: Int,
acks: Int = -1): Seq[String] = {
val values = (0 until numMessages).map(x => s"test-$x")
val records = values.map(v => new ProducerRecord[Array[Byte], Array[Byte]](topic, v.getBytes))
val intSerializer = new IntegerSerializer()
val records = values.zipWithIndex.map { case (v, i) =>
new ProducerRecord(topic, intSerializer.serialize(topic, i), v.getBytes)
}
produceMessages(servers, records, acks)
values
}
Expand Down