From 9d21b358a27ebfe3ca35da1098c934cd4dfed3e7 Mon Sep 17 00:00:00 2001 From: Jeff Kim Date: Wed, 5 Apr 2023 19:11:45 -0400 Subject: [PATCH 1/6] ignore unknown record types for coordinators --- .../kafka/common/protocol/MessageUtil.java | 34 ++++ .../group/GroupMetadataManager.scala | 149 ++++++++++-------- .../transaction/TransactionLog.scala | 74 +++++---- .../transaction/TransactionStateManager.scala | 22 +-- .../group/GroupMetadataManagerTest.scala | 21 ++- .../transaction/TransactionLogTest.scala | 9 +- .../TransactionStateManagerTest.scala | 40 ++++- 7 files changed, 237 insertions(+), 112 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/common/protocol/MessageUtil.java b/clients/src/main/java/org/apache/kafka/common/protocol/MessageUtil.java index 95bc01945ec7c..f753dcb0e5d40 100644 --- a/clients/src/main/java/org/apache/kafka/common/protocol/MessageUtil.java +++ b/clients/src/main/java/org/apache/kafka/common/protocol/MessageUtil.java @@ -231,4 +231,38 @@ public static byte[] toVersionPrefixedBytes(final short version, final Message m buffer.limit() == buffer.array().length) return buffer.array(); else return Utils.toArray(buffer); } + + // Should only be used for testing + public static byte[] messageWithUnknownVersion() { + return MessageUtil.toVersionPrefixedBytes(Short.MAX_VALUE, new Message() { + @Override + public short lowestSupportedVersion() { + return Short.MAX_VALUE; + } + + @Override + public short highestSupportedVersion() { + return Short.MAX_VALUE; + } + + @Override + public void addSize(MessageSizeAccumulator size, ObjectSerializationCache cache, short version) {} + + @Override + public void write(Writable writable, ObjectSerializationCache cache, short version) {} + + @Override + public void read(Readable readable, short version) {} + + @Override + public List unknownTaggedFields() { + return null; + } + + @Override + public Message duplicate() { + return null; + } + }); + } } diff --git a/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala b/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala index 2ea0d79662e73..18547330bb365 100644 --- a/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala +++ b/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala @@ -576,7 +576,10 @@ class GroupMetadataManager(brokerId: Int, } } - private def doLoadGroupsAndOffsets(topicPartition: TopicPartition, onGroupLoaded: GroupMetadata => Unit): Unit = { + // Visible for testing + private[group] def doLoadGroupsAndOffsets(topicPartition: TopicPartition, + onGroupLoaded: GroupMetadata => Unit): Unit = { + def logEndOffset: Long = replicaManager.getLogEndOffset(topicPartition).getOrElse(-1L) replicaManager.getLog(topicPartition) match { @@ -651,40 +654,42 @@ class GroupMetadataManager(brokerId: Int, if (batchBaseOffset.isEmpty) batchBaseOffset = Some(record.offset) GroupMetadataManager.readMessageKey(record.key) match { + case Some(key) => key match { + case offsetKey: OffsetKey => + if (isTxnOffsetCommit && !pendingOffsets.contains(batch.producerId)) + pendingOffsets.put(batch.producerId, mutable.Map[GroupTopicPartition, CommitRecordMetadataAndOffset]()) + + // load offset + val groupTopicPartition = offsetKey.key + if (!record.hasValue) { + if (isTxnOffsetCommit) + pendingOffsets(batch.producerId).remove(groupTopicPartition) + else + loadedOffsets.remove(groupTopicPartition) + } else { + val offsetAndMetadata = GroupMetadataManager.readOffsetMessageValue(record.value) + if (isTxnOffsetCommit) + pendingOffsets(batch.producerId).put(groupTopicPartition, CommitRecordMetadataAndOffset(batchBaseOffset, offsetAndMetadata)) + else + loadedOffsets.put(groupTopicPartition, CommitRecordMetadataAndOffset(batchBaseOffset, offsetAndMetadata)) + } + + case groupMetadataKey: GroupMetadataKey => + // load group metadata + val groupId = groupMetadataKey.key + val groupMetadata = GroupMetadataManager.readGroupMessageValue(groupId, record.value, time) + if (groupMetadata != null) { + removedGroups.remove(groupId) + loadedGroups.put(groupId, groupMetadata) + } else { + loadedGroups.remove(groupId) + removedGroups.add(groupId) + } + + case _ => // do nothing + } - case offsetKey: OffsetKey => - if (isTxnOffsetCommit && !pendingOffsets.contains(batch.producerId)) - pendingOffsets.put(batch.producerId, mutable.Map[GroupTopicPartition, CommitRecordMetadataAndOffset]()) - - // load offset - val groupTopicPartition = offsetKey.key - if (!record.hasValue) { - if (isTxnOffsetCommit) - pendingOffsets(batch.producerId).remove(groupTopicPartition) - else - loadedOffsets.remove(groupTopicPartition) - } else { - val offsetAndMetadata = GroupMetadataManager.readOffsetMessageValue(record.value) - if (isTxnOffsetCommit) - pendingOffsets(batch.producerId).put(groupTopicPartition, CommitRecordMetadataAndOffset(batchBaseOffset, offsetAndMetadata)) - else - loadedOffsets.put(groupTopicPartition, CommitRecordMetadataAndOffset(batchBaseOffset, offsetAndMetadata)) - } - - case groupMetadataKey: GroupMetadataKey => - // load group metadata - val groupId = groupMetadataKey.key - val groupMetadata = GroupMetadataManager.readGroupMessageValue(groupId, record.value, time) - if (groupMetadata != null) { - removedGroups.remove(groupId) - loadedGroups.put(groupId, groupMetadata) - } else { - loadedGroups.remove(groupId) - removedGroups.add(groupId) - } - - case unknownKey => - throw new IllegalStateException(s"Unexpected message key $unknownKey while loading offsets and group metadata") + case None => // ignore unknown keys } } } @@ -1041,7 +1046,7 @@ class GroupMetadataManager(brokerId: Int, * key version 2: group metadata * -> value version 0: [protocol_type, generation, protocol, leader, members] */ -object GroupMetadataManager { +object GroupMetadataManager extends Logging { // Metrics names val MetricsGroup: String = "group-coordinator-metrics" val LoadTimeSensor: String = "GroupPartitionLoadTime" @@ -1145,17 +1150,22 @@ object GroupMetadataManager { * @param buffer input byte-buffer * @return an OffsetKey or GroupMetadataKey object from the message */ - def readMessageKey(buffer: ByteBuffer): BaseKey = { + def readMessageKey(buffer: ByteBuffer): Option[BaseKey] = { val version = buffer.getShort if (version >= OffsetCommitKey.LOWEST_SUPPORTED_VERSION && version <= OffsetCommitKey.HIGHEST_SUPPORTED_VERSION) { // version 0 and 1 refer to offset val key = new OffsetCommitKey(new ByteBufferAccessor(buffer), version) - OffsetKey(version, GroupTopicPartition(key.group, new TopicPartition(key.topic, key.partition))) + Some(OffsetKey(version, GroupTopicPartition(key.group, new TopicPartition(key.topic, key.partition)))) } else if (version >= GroupMetadataKeyData.LOWEST_SUPPORTED_VERSION && version <= GroupMetadataKeyData.HIGHEST_SUPPORTED_VERSION) { // version 2 refers to group metadata val key = new GroupMetadataKeyData(new ByteBufferAccessor(buffer), version) - GroupMetadataKey(version, key.group) - } else throw new IllegalStateException(s"Unknown group metadata message version: $version") + Some(GroupMetadataKey(version, key.group)) + } else { + // Unknown versions may exist when a downgraded coordinator is reading records from the log. + warn(s"Found unknown message key version: $version." + + s" The downgraded coordinator will ignore this key and corresponding value.") + None + } } /** @@ -1229,17 +1239,21 @@ object GroupMetadataManager { Option(consumerRecord.key).map(key => GroupMetadataManager.readMessageKey(ByteBuffer.wrap(key))).foreach { // Only print if the message is an offset record. // We ignore the timestamp of the message because GroupMetadataMessage has its own timestamp. - case offsetKey: OffsetKey => - val groupTopicPartition = offsetKey.key - val value = consumerRecord.value - val formattedValue = - if (value == null) "NULL" - else GroupMetadataManager.readOffsetMessageValue(ByteBuffer.wrap(value)).toString - output.write(groupTopicPartition.toString.getBytes(StandardCharsets.UTF_8)) - output.write("::".getBytes(StandardCharsets.UTF_8)) - output.write(formattedValue.getBytes(StandardCharsets.UTF_8)) - output.write("\n".getBytes(StandardCharsets.UTF_8)) - case _ => // no-op + case Some(key) => key match { + case offsetKey: OffsetKey => + val groupTopicPartition = offsetKey.key + val value = consumerRecord.value + val formattedValue = + if (value == null) "NULL" + else GroupMetadataManager.readOffsetMessageValue(ByteBuffer.wrap(value)).toString + output.write(groupTopicPartition.toString.getBytes(StandardCharsets.UTF_8)) + output.write("::".getBytes(StandardCharsets.UTF_8)) + output.write(formattedValue.getBytes(StandardCharsets.UTF_8)) + output.write("\n".getBytes(StandardCharsets.UTF_8)) + case _ => // no-op + } + + case None => // no-op } } } @@ -1250,17 +1264,20 @@ object GroupMetadataManager { Option(consumerRecord.key).map(key => GroupMetadataManager.readMessageKey(ByteBuffer.wrap(key))).foreach { // Only print if the message is a group metadata record. // We ignore the timestamp of the message because GroupMetadataMessage has its own timestamp. - case groupMetadataKey: GroupMetadataKey => - val groupId = groupMetadataKey.key - val value = consumerRecord.value - val formattedValue = - if (value == null) "NULL" - else GroupMetadataManager.readGroupMessageValue(groupId, ByteBuffer.wrap(value), Time.SYSTEM).toString - output.write(groupId.getBytes(StandardCharsets.UTF_8)) - output.write("::".getBytes(StandardCharsets.UTF_8)) - output.write(formattedValue.getBytes(StandardCharsets.UTF_8)) - output.write("\n".getBytes(StandardCharsets.UTF_8)) - case _ => // no-op + case Some(key) => key match { + case groupMetadataKey: GroupMetadataKey => + val groupId = groupMetadataKey.key + val value = consumerRecord.value + val formattedValue = + if (value == null) "NULL" + else GroupMetadataManager.readGroupMessageValue(groupId, ByteBuffer.wrap(value), Time.SYSTEM).toString + output.write(groupId.getBytes(StandardCharsets.UTF_8)) + output.write("::".getBytes(StandardCharsets.UTF_8)) + output.write(formattedValue.getBytes(StandardCharsets.UTF_8)) + output.write("\n".getBytes(StandardCharsets.UTF_8)) + case _ => // no-op + } + case None => // no-op } } } @@ -1273,9 +1290,13 @@ object GroupMetadataManager { throw new KafkaException("Failed to decode message using offset topic decoder (message had a missing key)") } else { GroupMetadataManager.readMessageKey(record.key) match { - case offsetKey: OffsetKey => parseOffsets(offsetKey, record.value) - case groupMetadataKey: GroupMetadataKey => parseGroupMetadata(groupMetadataKey, record.value) - case _ => throw new KafkaException("Failed to decode message using offset topic decoder (message had an invalid key)") + case Some(key) => key match { + case offsetKey: OffsetKey => parseOffsets(offsetKey, record.value) + case groupMetadataKey: GroupMetadataKey => parseGroupMetadata(groupMetadataKey, record.value) + case _ => throw new KafkaException("Failed to decode message using offset topic decoder (message had an invalid key)") + } + case None => + (Some(""), Some("")) } } } diff --git a/core/src/main/scala/kafka/coordinator/transaction/TransactionLog.scala b/core/src/main/scala/kafka/coordinator/transaction/TransactionLog.scala index cb501f774fd9d..ed0dbf07086b3 100644 --- a/core/src/main/scala/kafka/coordinator/transaction/TransactionLog.scala +++ b/core/src/main/scala/kafka/coordinator/transaction/TransactionLog.scala @@ -19,8 +19,8 @@ package kafka.coordinator.transaction import java.io.PrintStream import java.nio.ByteBuffer import java.nio.charset.StandardCharsets - import kafka.internals.generated.{TransactionLogKey, TransactionLogValue} +import kafka.utils.Logging import org.apache.kafka.clients.consumer.ConsumerRecord import org.apache.kafka.common.protocol.{ByteBufferAccessor, MessageUtil} import org.apache.kafka.common.record.{CompressionType, Record, RecordBatch} @@ -37,7 +37,7 @@ import scala.jdk.CollectionConverters._ * key version 0: [transactionalId] * -> value version 0: [producer_id, producer_epoch, expire_timestamp, status, [topic, [partition] ], timestamp] */ -object TransactionLog { +object TransactionLog extends Logging { // log-level config default values and enforced values val DefaultNumPartitions: Int = 50 @@ -98,15 +98,19 @@ object TransactionLog { * * @return the key */ - def readTxnRecordKey(buffer: ByteBuffer): TxnKey = { + def readTxnRecordKey(buffer: ByteBuffer): Option[TxnKey] = { val version = buffer.getShort if (version >= TransactionLogKey.LOWEST_SUPPORTED_VERSION && version <= TransactionLogKey.HIGHEST_SUPPORTED_VERSION) { val value = new TransactionLogKey(new ByteBufferAccessor(buffer), version) - TxnKey( + Some(TxnKey( version = version, transactionalId = value.transactionalId - ) - } else throw new IllegalStateException(s"Unknown version $version from the transaction log message") + )) + } else { + warn(s"Unknown version $version from the transaction log message." + + s" The downgraded coordinator will ignore this key and corresponding value.") + None + } } /** @@ -148,17 +152,20 @@ object TransactionLog { // Formatter for use with tools to read transaction log messages class TransactionLogMessageFormatter extends MessageFormatter { def writeTo(consumerRecord: ConsumerRecord[Array[Byte], Array[Byte]], output: PrintStream): Unit = { - Option(consumerRecord.key).map(key => readTxnRecordKey(ByteBuffer.wrap(key))).foreach { txnKey => - val transactionalId = txnKey.transactionalId - val value = consumerRecord.value - val producerIdMetadata = if (value == null) - None - else - readTxnRecordValue(transactionalId, ByteBuffer.wrap(value)) - output.write(transactionalId.getBytes(StandardCharsets.UTF_8)) - output.write("::".getBytes(StandardCharsets.UTF_8)) - output.write(producerIdMetadata.getOrElse("NULL").toString.getBytes(StandardCharsets.UTF_8)) - output.write("\n".getBytes(StandardCharsets.UTF_8)) + Option(consumerRecord.key).map(key => readTxnRecordKey(ByteBuffer.wrap(key))).foreach { + case Some(txnKey) => + val transactionalId = txnKey.transactionalId + val value = consumerRecord.value + val producerIdMetadata = if (value == null) + None + else + readTxnRecordValue(transactionalId, ByteBuffer.wrap(value)) + output.write(transactionalId.getBytes(StandardCharsets.UTF_8)) + output.write("::".getBytes(StandardCharsets.UTF_8)) + output.write(producerIdMetadata.getOrElse("NULL").toString.getBytes(StandardCharsets.UTF_8)) + output.write("\n".getBytes(StandardCharsets.UTF_8)) + + case None => // Only print if this message is a transaction record } } } @@ -167,21 +174,26 @@ object TransactionLog { * Exposed for printing records using [[kafka.tools.DumpLogSegments]] */ def formatRecordKeyAndValue(record: Record): (Option[String], Option[String]) = { - val txnKey = TransactionLog.readTxnRecordKey(record.key) - val keyString = s"transaction_metadata::transactionalId=${txnKey.transactionalId}" - - val valueString = TransactionLog.readTxnRecordValue(txnKey.transactionalId, record.value) match { - case None => "" - - case Some(txnMetadata) => s"producerId:${txnMetadata.producerId}," + - s"producerEpoch:${txnMetadata.producerEpoch}," + - s"state=${txnMetadata.state}," + - s"partitions=${txnMetadata.topicPartitions.mkString("[", ",", "]")}," + - s"txnLastUpdateTimestamp=${txnMetadata.txnLastUpdateTimestamp}," + - s"txnTimeoutMs=${txnMetadata.txnTimeoutMs}" - } + TransactionLog.readTxnRecordKey(record.key) match { + case Some(txnKey) => + val keyString = s"transaction_metadata::transactionalId=${txnKey.transactionalId}" + + val valueString = TransactionLog.readTxnRecordValue(txnKey.transactionalId, record.value) match { + case None => "" - (Some(keyString), Some(valueString)) + case Some(txnMetadata) => s"producerId:${txnMetadata.producerId}," + + s"producerEpoch:${txnMetadata.producerEpoch}," + + s"state=${txnMetadata.state}," + + s"partitions=${txnMetadata.topicPartitions.mkString("[", ",", "]")}," + + s"txnLastUpdateTimestamp=${txnMetadata.txnLastUpdateTimestamp}," + + s"txnTimeoutMs=${txnMetadata.txnTimeoutMs}" + } + + (Some(keyString), Some(valueString)) + + case None => + (Some(""), Some("")) + } } } diff --git a/core/src/main/scala/kafka/coordinator/transaction/TransactionStateManager.scala b/core/src/main/scala/kafka/coordinator/transaction/TransactionStateManager.scala index f561559fa635d..6395470b4c7ee 100644 --- a/core/src/main/scala/kafka/coordinator/transaction/TransactionStateManager.scala +++ b/core/src/main/scala/kafka/coordinator/transaction/TransactionStateManager.scala @@ -467,16 +467,20 @@ class TransactionStateManager(brokerId: Int, memRecords.batches.forEach { batch => for (record <- batch.asScala) { require(record.hasKey, "Transaction state log's key should not be null") - val txnKey = TransactionLog.readTxnRecordKey(record.key) - // load transaction metadata along with transaction state - val transactionalId = txnKey.transactionalId - TransactionLog.readTxnRecordValue(transactionalId, record.value) match { - case None => - loadedTransactions.remove(transactionalId) - case Some(txnMetadata) => - loadedTransactions.put(transactionalId, txnMetadata) + TransactionLog.readTxnRecordKey(record.key) match { + case Some(txnKey) => + // load transaction metadata along with transaction state + val transactionalId = txnKey.transactionalId + TransactionLog.readTxnRecordValue(transactionalId, record.value) match { + case None => + loadedTransactions.remove(transactionalId) + case Some(txnMetadata) => + loadedTransactions.put(transactionalId, txnMetadata) + } + currOffset = batch.nextOffset + + case None => // ignore unknown keys } - currOffset = batch.nextOffset } } } diff --git a/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala b/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala index 333aeac5af711..ce5fc780d532b 100644 --- a/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala @@ -35,7 +35,7 @@ import org.apache.kafka.clients.consumer.internals.ConsumerProtocol import org.apache.kafka.common.{TopicIdPartition, TopicPartition, Uuid} import org.apache.kafka.common.internals.Topic import org.apache.kafka.common.metrics.{JmxReporter, KafkaMetricsContext, Metrics => kMetrics} -import org.apache.kafka.common.protocol.Errors +import org.apache.kafka.common.protocol.{Errors, MessageUtil} import org.apache.kafka.common.record._ import org.apache.kafka.common.requests.OffsetFetchResponse import org.apache.kafka.common.requests.ProduceResponse.PartitionResponse @@ -640,8 +640,13 @@ class GroupMetadataManagerTest { val offsetCommitRecords = createCommittedOffsetRecords(committedOffsets) val memberId = "98098230493" val groupMetadataRecord = buildStableGroupRecordWithMember(generation, protocolType, protocol, memberId) + + // Should ignore unknown record + val unknownMessage = MessageUtil.messageWithUnknownVersion() + val unknownRecord = new SimpleRecord(unknownMessage, unknownMessage) + val records = MemoryRecords.withRecords(startOffset, CompressionType.NONE, - (offsetCommitRecords ++ Seq(groupMetadataRecord)).toArray: _*) + (offsetCommitRecords ++ Seq(unknownRecord) ++ Seq(groupMetadataRecord)).toArray: _*) expectGroupMetadataLoad(groupMetadataTopicPartition, startOffset, records) @@ -1727,7 +1732,7 @@ class GroupMetadataManagerTest { assertFalse(metadataTombstone.hasValue) assertTrue(metadataTombstone.timestamp > 0) - val groupKey = GroupMetadataManager.readMessageKey(metadataTombstone.key).asInstanceOf[GroupMetadataKey] + val groupKey = GroupMetadataManager.readMessageKey(metadataTombstone.key).get.asInstanceOf[GroupMetadataKey] assertEquals(groupId, groupKey.key) // the full group should be gone since all offsets were removed @@ -1771,7 +1776,7 @@ class GroupMetadataManagerTest { assertFalse(metadataTombstone.hasValue) assertTrue(metadataTombstone.timestamp > 0) - val groupKey = GroupMetadataManager.readMessageKey(metadataTombstone.key).asInstanceOf[GroupMetadataKey] + val groupKey = GroupMetadataManager.readMessageKey(metadataTombstone.key).get.asInstanceOf[GroupMetadataKey] assertEquals(groupId, groupKey.key) // the full group should be gone since all offsets were removed @@ -1832,7 +1837,7 @@ class GroupMetadataManagerTest { records.foreach { message => assertTrue(message.hasKey) assertFalse(message.hasValue) - val offsetKey = GroupMetadataManager.readMessageKey(message.key).asInstanceOf[OffsetKey] + val offsetKey = GroupMetadataManager.readMessageKey(message.key).get.asInstanceOf[OffsetKey] assertEquals(groupId, offsetKey.key.group) assertEquals("foo", offsetKey.key.topicPartition.topic) } @@ -2762,4 +2767,10 @@ class GroupMetadataManagerTest { assertTrue(partitionLoadTime("partition-load-time-max") >= diff) assertTrue(partitionLoadTime("partition-load-time-avg") >= diff) } + + @Test + def testIgnoreUnknownMessageKeyVersion(): Unit = { + GroupMetadataManager.readMessageKey(ByteBuffer.wrap(MessageUtil.messageWithUnknownVersion())) + } + } diff --git a/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionLogTest.scala b/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionLogTest.scala index 32e17d88a7b1e..1e6bbe94a1022 100644 --- a/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionLogTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionLogTest.scala @@ -19,10 +19,12 @@ package kafka.coordinator.transaction import kafka.utils.TestUtils import org.apache.kafka.common.TopicPartition +import org.apache.kafka.common.protocol.MessageUtil import org.apache.kafka.common.record.{CompressionType, MemoryRecords, SimpleRecord} import org.junit.jupiter.api.Assertions.{assertEquals, assertThrows} import org.junit.jupiter.api.Test +import java.nio.ByteBuffer import scala.jdk.CollectionConverters._ class TransactionLogTest { @@ -81,7 +83,7 @@ class TransactionLogTest { var count = 0 for (record <- records.records.asScala) { - val txnKey = TransactionLog.readTxnRecordKey(record.key) + val txnKey = TransactionLog.readTxnRecordKey(record.key).get val transactionalId = txnKey.transactionalId val txnMetadata = TransactionLog.readTxnRecordValue(transactionalId, record.value).get @@ -135,4 +137,9 @@ class TransactionLogTest { assertEquals(Some(""), valueStringOpt) } + @Test + def testReadUnknownMessageKeyVersion(): Unit = { + TransactionLog.readTxnRecordKey(ByteBuffer.wrap(MessageUtil.messageWithUnknownVersion())) + } + } diff --git a/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionStateManagerTest.scala b/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionStateManagerTest.scala index 21f27feb2a1d0..dcf31d5c48662 100644 --- a/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionStateManagerTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionStateManagerTest.scala @@ -28,7 +28,7 @@ import kafka.zk.KafkaZkClient import org.apache.kafka.common.TopicPartition import org.apache.kafka.common.internals.Topic.TRANSACTION_STATE_TOPIC_NAME import org.apache.kafka.common.metrics.{JmxReporter, KafkaMetricsContext, Metrics} -import org.apache.kafka.common.protocol.Errors +import org.apache.kafka.common.protocol.{Errors, MessageUtil} import org.apache.kafka.common.record._ import org.apache.kafka.common.requests.ProduceResponse.PartitionResponse import org.apache.kafka.common.requests.TransactionResult @@ -758,7 +758,7 @@ class TransactionStateManagerTest { appendedRecords.values.foreach { batches => batches.foreach { records => records.records.forEach { record => - val transactionalId = TransactionLog.readTxnRecordKey(record.key).transactionalId + val transactionalId = TransactionLog.readTxnRecordKey(record.key).get.transactionalId assertNull(record.value) expiredTransactionalIds += transactionalId assertEquals(Right(None), transactionManager.getTransactionState(transactionalId)) @@ -1085,4 +1085,40 @@ class TransactionStateManagerTest { assertTrue(partitionLoadTime("partition-load-time-max") >= 0) assertTrue(partitionLoadTime( "partition-load-time-avg") >= 0) } + + + @Test + def testIgnoreUnknownRecordType(): Unit = { + txnMetadata1.state = PrepareCommit + txnMetadata1.addPartitions(Set[TopicPartition](new TopicPartition("topic1", 0), + new TopicPartition("topic1", 1))) + + txnRecords += new SimpleRecord(txnMessageKeyBytes1, TransactionLog.valueToBytes(txnMetadata1.prepareNoTransit())) + val startOffset = 0L + + val unknownMessage = MessageUtil.messageWithUnknownVersion() + val unknownRecord = new SimpleRecord(unknownMessage, unknownMessage) + + val records = MemoryRecords.withRecords(startOffset, CompressionType.NONE, + (Seq(unknownRecord) ++ txnRecords).toArray: _*) + + prepareTxnLog(topicPartition, 0, records) + + transactionManager.loadTransactionsForTxnTopicPartition(partitionId, coordinatorEpoch = 1, (_, _, _, _) => ()) + assertEquals(0, transactionManager.loadingPartitions.size) + assertTrue(transactionManager.transactionMetadataCache.contains(partitionId)) + val txnMetadataPool = transactionManager.transactionMetadataCache(partitionId).metadataPerTransactionalId + assertFalse(txnMetadataPool.isEmpty) + assertTrue(txnMetadataPool.contains(transactionalId1)) + val txnMetadata = txnMetadataPool.get(transactionalId1) + assertEquals(txnMetadata1.transactionalId, txnMetadata.transactionalId) + assertEquals(txnMetadata1.producerId, txnMetadata.producerId) + assertEquals(txnMetadata1.lastProducerId, txnMetadata.lastProducerId) + assertEquals(txnMetadata1.producerEpoch, txnMetadata.producerEpoch) + assertEquals(txnMetadata1.lastProducerEpoch, txnMetadata.lastProducerEpoch) + assertEquals(txnMetadata1.txnTimeoutMs, txnMetadata.txnTimeoutMs) + assertEquals(txnMetadata1.state, txnMetadata.state) + assertEquals(txnMetadata1.topicPartitions, txnMetadata.topicPartitions) + assertEquals(1, transactionManager.transactionMetadataCache(partitionId).coordinatorEpoch) + } } From dd57136746dab829782d3d363fd5f199877df520 Mon Sep 17 00:00:00 2001 From: Jeff Kim Date: Thu, 13 Apr 2023 10:32:05 -0400 Subject: [PATCH 2/6] address comments --- .../kafka/common/protocol/MessageUtil.java | 34 ----- .../group/GroupMetadataManager.scala | 138 +++++++++--------- .../transaction/TransactionLog.scala | 34 +++-- .../transaction/TransactionStateManager.scala | 7 +- .../group/GroupMetadataManagerTest.scala | 10 +- .../transaction/TransactionLogTest.scala | 7 +- 6 files changed, 107 insertions(+), 123 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/common/protocol/MessageUtil.java b/clients/src/main/java/org/apache/kafka/common/protocol/MessageUtil.java index f753dcb0e5d40..95bc01945ec7c 100644 --- a/clients/src/main/java/org/apache/kafka/common/protocol/MessageUtil.java +++ b/clients/src/main/java/org/apache/kafka/common/protocol/MessageUtil.java @@ -231,38 +231,4 @@ public static byte[] toVersionPrefixedBytes(final short version, final Message m buffer.limit() == buffer.array().length) return buffer.array(); else return Utils.toArray(buffer); } - - // Should only be used for testing - public static byte[] messageWithUnknownVersion() { - return MessageUtil.toVersionPrefixedBytes(Short.MAX_VALUE, new Message() { - @Override - public short lowestSupportedVersion() { - return Short.MAX_VALUE; - } - - @Override - public short highestSupportedVersion() { - return Short.MAX_VALUE; - } - - @Override - public void addSize(MessageSizeAccumulator size, ObjectSerializationCache cache, short version) {} - - @Override - public void write(Writable writable, ObjectSerializationCache cache, short version) {} - - @Override - public void read(Readable readable, short version) {} - - @Override - public List unknownTaggedFields() { - return null; - } - - @Override - public Message duplicate() { - return null; - } - }); - } } diff --git a/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala b/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala index 18547330bb365..3174631c47f15 100644 --- a/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala +++ b/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala @@ -654,42 +654,41 @@ class GroupMetadataManager(brokerId: Int, if (batchBaseOffset.isEmpty) batchBaseOffset = Some(record.offset) GroupMetadataManager.readMessageKey(record.key) match { - case Some(key) => key match { - case offsetKey: OffsetKey => - if (isTxnOffsetCommit && !pendingOffsets.contains(batch.producerId)) - pendingOffsets.put(batch.producerId, mutable.Map[GroupTopicPartition, CommitRecordMetadataAndOffset]()) - - // load offset - val groupTopicPartition = offsetKey.key - if (!record.hasValue) { - if (isTxnOffsetCommit) - pendingOffsets(batch.producerId).remove(groupTopicPartition) - else - loadedOffsets.remove(groupTopicPartition) - } else { - val offsetAndMetadata = GroupMetadataManager.readOffsetMessageValue(record.value) - if (isTxnOffsetCommit) - pendingOffsets(batch.producerId).put(groupTopicPartition, CommitRecordMetadataAndOffset(batchBaseOffset, offsetAndMetadata)) - else - loadedOffsets.put(groupTopicPartition, CommitRecordMetadataAndOffset(batchBaseOffset, offsetAndMetadata)) - } - - case groupMetadataKey: GroupMetadataKey => - // load group metadata - val groupId = groupMetadataKey.key - val groupMetadata = GroupMetadataManager.readGroupMessageValue(groupId, record.value, time) - if (groupMetadata != null) { - removedGroups.remove(groupId) - loadedGroups.put(groupId, groupMetadata) - } else { - loadedGroups.remove(groupId) - removedGroups.add(groupId) - } - - case _ => // do nothing - } + case offsetKey: OffsetKey => + if (isTxnOffsetCommit && !pendingOffsets.contains(batch.producerId)) + pendingOffsets.put(batch.producerId, mutable.Map[GroupTopicPartition, CommitRecordMetadataAndOffset]()) + + // load offset + val groupTopicPartition = offsetKey.key + if (!record.hasValue) { + if (isTxnOffsetCommit) + pendingOffsets(batch.producerId).remove(groupTopicPartition) + else + loadedOffsets.remove(groupTopicPartition) + } else { + val offsetAndMetadata = GroupMetadataManager.readOffsetMessageValue(record.value) + if (isTxnOffsetCommit) + pendingOffsets(batch.producerId).put(groupTopicPartition, CommitRecordMetadataAndOffset(batchBaseOffset, offsetAndMetadata)) + else + loadedOffsets.put(groupTopicPartition, CommitRecordMetadataAndOffset(batchBaseOffset, offsetAndMetadata)) + } + + case groupMetadataKey: GroupMetadataKey => + // load group metadata + val groupId = groupMetadataKey.key + val groupMetadata = GroupMetadataManager.readGroupMessageValue(groupId, record.value, time) + if (groupMetadata != null) { + removedGroups.remove(groupId) + loadedGroups.put(groupId, groupMetadata) + } else { + loadedGroups.remove(groupId) + removedGroups.add(groupId) + } - case None => // ignore unknown keys + case _: UnknownKey => // do nothing + + case unexpectedKey => + throw new IllegalStateException(s"Unexpected message key $unexpectedKey while loading offsets and group metadata") } } } @@ -1150,21 +1149,21 @@ object GroupMetadataManager extends Logging { * @param buffer input byte-buffer * @return an OffsetKey or GroupMetadataKey object from the message */ - def readMessageKey(buffer: ByteBuffer): Option[BaseKey] = { + def readMessageKey(buffer: ByteBuffer): BaseKey = { val version = buffer.getShort if (version >= OffsetCommitKey.LOWEST_SUPPORTED_VERSION && version <= OffsetCommitKey.HIGHEST_SUPPORTED_VERSION) { // version 0 and 1 refer to offset val key = new OffsetCommitKey(new ByteBufferAccessor(buffer), version) - Some(OffsetKey(version, GroupTopicPartition(key.group, new TopicPartition(key.topic, key.partition)))) + OffsetKey(version, GroupTopicPartition(key.group, new TopicPartition(key.topic, key.partition))) } else if (version >= GroupMetadataKeyData.LOWEST_SUPPORTED_VERSION && version <= GroupMetadataKeyData.HIGHEST_SUPPORTED_VERSION) { // version 2 refers to group metadata val key = new GroupMetadataKeyData(new ByteBufferAccessor(buffer), version) - Some(GroupMetadataKey(version, key.group)) + GroupMetadataKey(version, key.group) } else { // Unknown versions may exist when a downgraded coordinator is reading records from the log. warn(s"Found unknown message key version: $version." + s" The downgraded coordinator will ignore this key and corresponding value.") - None + UnknownKey(version) } } @@ -1239,21 +1238,17 @@ object GroupMetadataManager extends Logging { Option(consumerRecord.key).map(key => GroupMetadataManager.readMessageKey(ByteBuffer.wrap(key))).foreach { // Only print if the message is an offset record. // We ignore the timestamp of the message because GroupMetadataMessage has its own timestamp. - case Some(key) => key match { - case offsetKey: OffsetKey => - val groupTopicPartition = offsetKey.key - val value = consumerRecord.value - val formattedValue = - if (value == null) "NULL" - else GroupMetadataManager.readOffsetMessageValue(ByteBuffer.wrap(value)).toString - output.write(groupTopicPartition.toString.getBytes(StandardCharsets.UTF_8)) - output.write("::".getBytes(StandardCharsets.UTF_8)) - output.write(formattedValue.getBytes(StandardCharsets.UTF_8)) - output.write("\n".getBytes(StandardCharsets.UTF_8)) - case _ => // no-op - } - - case None => // no-op + case offsetKey: OffsetKey => + val groupTopicPartition = offsetKey.key + val value = consumerRecord.value + val formattedValue = + if (value == null) "NULL" + else GroupMetadataManager.readOffsetMessageValue(ByteBuffer.wrap(value)).toString + output.write(groupTopicPartition.toString.getBytes(StandardCharsets.UTF_8)) + output.write("::".getBytes(StandardCharsets.UTF_8)) + output.write(formattedValue.getBytes(StandardCharsets.UTF_8)) + output.write("\n".getBytes(StandardCharsets.UTF_8)) + case _ => // no-op } } } @@ -1264,20 +1259,17 @@ object GroupMetadataManager extends Logging { Option(consumerRecord.key).map(key => GroupMetadataManager.readMessageKey(ByteBuffer.wrap(key))).foreach { // Only print if the message is a group metadata record. // We ignore the timestamp of the message because GroupMetadataMessage has its own timestamp. - case Some(key) => key match { - case groupMetadataKey: GroupMetadataKey => - val groupId = groupMetadataKey.key - val value = consumerRecord.value - val formattedValue = - if (value == null) "NULL" - else GroupMetadataManager.readGroupMessageValue(groupId, ByteBuffer.wrap(value), Time.SYSTEM).toString - output.write(groupId.getBytes(StandardCharsets.UTF_8)) - output.write("::".getBytes(StandardCharsets.UTF_8)) - output.write(formattedValue.getBytes(StandardCharsets.UTF_8)) - output.write("\n".getBytes(StandardCharsets.UTF_8)) - case _ => // no-op - } - case None => // no-op + case groupMetadataKey: GroupMetadataKey => + val groupId = groupMetadataKey.key + val value = consumerRecord.value + val formattedValue = + if (value == null) "NULL" + else GroupMetadataManager.readGroupMessageValue(groupId, ByteBuffer.wrap(value), Time.SYSTEM).toString + output.write(groupId.getBytes(StandardCharsets.UTF_8)) + output.write("::".getBytes(StandardCharsets.UTF_8)) + output.write(formattedValue.getBytes(StandardCharsets.UTF_8)) + output.write("\n".getBytes(StandardCharsets.UTF_8)) + case _ => // no-op } } } @@ -1290,13 +1282,10 @@ object GroupMetadataManager extends Logging { throw new KafkaException("Failed to decode message using offset topic decoder (message had a missing key)") } else { GroupMetadataManager.readMessageKey(record.key) match { - case Some(key) => key match { case offsetKey: OffsetKey => parseOffsets(offsetKey, record.value) case groupMetadataKey: GroupMetadataKey => parseGroupMetadata(groupMetadataKey, record.value) + case _: UnknownKey => (Some(""), Some("")) case _ => throw new KafkaException("Failed to decode message using offset topic decoder (message had an invalid key)") - } - case None => - (Some(""), Some("")) } } } @@ -1389,3 +1378,8 @@ case class GroupMetadataKey(version: Short, key: String) extends BaseKey { override def toString: String = key } +case class UnknownKey(version: Short, key: String = null) extends BaseKey { + + override def toString: String = key +} + diff --git a/core/src/main/scala/kafka/coordinator/transaction/TransactionLog.scala b/core/src/main/scala/kafka/coordinator/transaction/TransactionLog.scala index ed0dbf07086b3..4623bb8a1a37b 100644 --- a/core/src/main/scala/kafka/coordinator/transaction/TransactionLog.scala +++ b/core/src/main/scala/kafka/coordinator/transaction/TransactionLog.scala @@ -98,18 +98,18 @@ object TransactionLog extends Logging { * * @return the key */ - def readTxnRecordKey(buffer: ByteBuffer): Option[TxnKey] = { + def readTxnRecordKey(buffer: ByteBuffer): BaseKey = { val version = buffer.getShort if (version >= TransactionLogKey.LOWEST_SUPPORTED_VERSION && version <= TransactionLogKey.HIGHEST_SUPPORTED_VERSION) { val value = new TransactionLogKey(new ByteBufferAccessor(buffer), version) - Some(TxnKey( + TxnKey( version = version, transactionalId = value.transactionalId - )) + ) } else { warn(s"Unknown version $version from the transaction log message." + s" The downgraded coordinator will ignore this key and corresponding value.") - None + UnknownKey(version) } } @@ -153,7 +153,7 @@ object TransactionLog extends Logging { class TransactionLogMessageFormatter extends MessageFormatter { def writeTo(consumerRecord: ConsumerRecord[Array[Byte], Array[Byte]], output: PrintStream): Unit = { Option(consumerRecord.key).map(key => readTxnRecordKey(ByteBuffer.wrap(key))).foreach { - case Some(txnKey) => + case txnKey: TxnKey => val transactionalId = txnKey.transactionalId val value = consumerRecord.value val producerIdMetadata = if (value == null) @@ -165,7 +165,10 @@ object TransactionLog extends Logging { output.write(producerIdMetadata.getOrElse("NULL").toString.getBytes(StandardCharsets.UTF_8)) output.write("\n".getBytes(StandardCharsets.UTF_8)) - case None => // Only print if this message is a transaction record + case _: UnknownKey => // Only print if this message is a transaction record + + case unexpectedKey => + throw new IllegalStateException(s"Found unexpected key $unexpectedKey while reading transaction log.") } } } @@ -175,7 +178,7 @@ object TransactionLog extends Logging { */ def formatRecordKeyAndValue(record: Record): (Option[String], Option[String]) = { TransactionLog.readTxnRecordKey(record.key) match { - case Some(txnKey) => + case txnKey: TxnKey => val keyString = s"transaction_metadata::transactionalId=${txnKey.transactionalId}" val valueString = TransactionLog.readTxnRecordValue(txnKey.transactionalId, record.value) match { @@ -191,13 +194,26 @@ object TransactionLog extends Logging { (Some(keyString), Some(valueString)) - case None => + case _: UnknownKey => (Some(""), Some("")) + + case unexpectedKey => + throw new IllegalStateException(s"Found unexpected key $unexpectedKey while formatting transaction log.") } } } -case class TxnKey(version: Short, transactionalId: String) { +trait BaseKey{ + def version: Short + def transactionalId: String +} + +case class TxnKey(version: Short, transactionalId: String) extends BaseKey { override def toString: String = transactionalId } + +case class UnknownKey(version: Short, transactionalId: String = null) extends BaseKey { + override def toString: String = transactionalId +} + diff --git a/core/src/main/scala/kafka/coordinator/transaction/TransactionStateManager.scala b/core/src/main/scala/kafka/coordinator/transaction/TransactionStateManager.scala index 6395470b4c7ee..85b0e472b4eda 100644 --- a/core/src/main/scala/kafka/coordinator/transaction/TransactionStateManager.scala +++ b/core/src/main/scala/kafka/coordinator/transaction/TransactionStateManager.scala @@ -468,7 +468,7 @@ class TransactionStateManager(brokerId: Int, for (record <- batch.asScala) { require(record.hasKey, "Transaction state log's key should not be null") TransactionLog.readTxnRecordKey(record.key) match { - case Some(txnKey) => + case txnKey: TxnKey => // load transaction metadata along with transaction state val transactionalId = txnKey.transactionalId TransactionLog.readTxnRecordValue(transactionalId, record.value) match { @@ -479,7 +479,10 @@ class TransactionStateManager(brokerId: Int, } currOffset = batch.nextOffset - case None => // ignore unknown keys + case _: UnknownKey => // ignore unknown keys + + case unexpectedKey => + throw new IllegalStateException(s"Found unexpected key $unexpectedKey while reading transaction log.") } } } diff --git a/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala b/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala index ce5fc780d532b..2329e841e3dce 100644 --- a/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala @@ -1732,7 +1732,7 @@ class GroupMetadataManagerTest { assertFalse(metadataTombstone.hasValue) assertTrue(metadataTombstone.timestamp > 0) - val groupKey = GroupMetadataManager.readMessageKey(metadataTombstone.key).get.asInstanceOf[GroupMetadataKey] + val groupKey = GroupMetadataManager.readMessageKey(metadataTombstone.key).asInstanceOf[GroupMetadataKey] assertEquals(groupId, groupKey.key) // the full group should be gone since all offsets were removed @@ -1776,7 +1776,7 @@ class GroupMetadataManagerTest { assertFalse(metadataTombstone.hasValue) assertTrue(metadataTombstone.timestamp > 0) - val groupKey = GroupMetadataManager.readMessageKey(metadataTombstone.key).get.asInstanceOf[GroupMetadataKey] + val groupKey = GroupMetadataManager.readMessageKey(metadataTombstone.key).asInstanceOf[GroupMetadataKey] assertEquals(groupId, groupKey.key) // the full group should be gone since all offsets were removed @@ -1837,7 +1837,7 @@ class GroupMetadataManagerTest { records.foreach { message => assertTrue(message.hasKey) assertFalse(message.hasValue) - val offsetKey = GroupMetadataManager.readMessageKey(message.key).get.asInstanceOf[OffsetKey] + val offsetKey = GroupMetadataManager.readMessageKey(message.key).asInstanceOf[OffsetKey] assertEquals(groupId, offsetKey.key.group) assertEquals("foo", offsetKey.key.topicPartition.topic) } @@ -2770,7 +2770,9 @@ class GroupMetadataManagerTest { @Test def testIgnoreUnknownMessageKeyVersion(): Unit = { - GroupMetadataManager.readMessageKey(ByteBuffer.wrap(MessageUtil.messageWithUnknownVersion())) + val record = new org.apache.kafka.coordinator.group.generated.GroupMetadataKey() + val unknownRecord = MessageUtil.toVersionPrefixedBytes(Short.MaxValue, record) + GroupMetadataManager.readMessageKey(ByteBuffer.wrap(unknownRecord)) } } diff --git a/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionLogTest.scala b/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionLogTest.scala index 1e6bbe94a1022..d8e72fe838403 100644 --- a/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionLogTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionLogTest.scala @@ -17,6 +17,7 @@ package kafka.coordinator.transaction +import kafka.internals.generated.TransactionLogKey import kafka.utils.TestUtils import org.apache.kafka.common.TopicPartition import org.apache.kafka.common.protocol.MessageUtil @@ -83,7 +84,7 @@ class TransactionLogTest { var count = 0 for (record <- records.records.asScala) { - val txnKey = TransactionLog.readTxnRecordKey(record.key).get + val txnKey = TransactionLog.readTxnRecordKey(record.key) val transactionalId = txnKey.transactionalId val txnMetadata = TransactionLog.readTxnRecordValue(transactionalId, record.value).get @@ -139,7 +140,9 @@ class TransactionLogTest { @Test def testReadUnknownMessageKeyVersion(): Unit = { - TransactionLog.readTxnRecordKey(ByteBuffer.wrap(MessageUtil.messageWithUnknownVersion())) + val record = new TransactionLogKey() + val unknownRecord = MessageUtil.toVersionPrefixedBytes(Short.MaxValue, record) + TransactionLog.readTxnRecordKey(ByteBuffer.wrap(unknownRecord)) } } From 3ae729cc893562b45d86f2be05059d11c9cef93f Mon Sep 17 00:00:00 2001 From: Jeff Kim Date: Thu, 13 Apr 2023 15:15:07 -0400 Subject: [PATCH 3/6] address comments --- .../group/GroupMetadataManager.scala | 20 +++++++++---------- .../transaction/TransactionLog.scala | 8 +++----- .../transaction/TransactionStateManager.scala | 5 ++++- .../group/GroupMetadataManagerTest.scala | 6 ++++-- .../transaction/TransactionLogTest.scala | 3 ++- .../TransactionStateManagerTest.scala | 7 +++++-- 6 files changed, 27 insertions(+), 22 deletions(-) diff --git a/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala b/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala index 3174631c47f15..4ce684aa6b1ca 100644 --- a/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala +++ b/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala @@ -685,7 +685,11 @@ class GroupMetadataManager(brokerId: Int, removedGroups.add(groupId) } - case _: UnknownKey => // do nothing + case unknownKey: UnknownKey => + // Unknown versions may exist when a downgraded coordinator is reading records from the log. + warn(s"Unknown message key with version ${unknownKey.version}" + + s" while loading offsets and group metadata. Ignoring it. " + + s"It could be a left over from an aborted upgrade.") case unexpectedKey => throw new IllegalStateException(s"Unexpected message key $unexpectedKey while loading offsets and group metadata") @@ -1045,7 +1049,7 @@ class GroupMetadataManager(brokerId: Int, * key version 2: group metadata * -> value version 0: [protocol_type, generation, protocol, leader, members] */ -object GroupMetadataManager extends Logging { +object GroupMetadataManager { // Metrics names val MetricsGroup: String = "group-coordinator-metrics" val LoadTimeSensor: String = "GroupPartitionLoadTime" @@ -1160,9 +1164,6 @@ object GroupMetadataManager extends Logging { val key = new GroupMetadataKeyData(new ByteBufferAccessor(buffer), version) GroupMetadataKey(version, key.group) } else { - // Unknown versions may exist when a downgraded coordinator is reading records from the log. - warn(s"Found unknown message key version: $version." + - s" The downgraded coordinator will ignore this key and corresponding value.") UnknownKey(version) } } @@ -1284,7 +1285,7 @@ object GroupMetadataManager extends Logging { GroupMetadataManager.readMessageKey(record.key) match { case offsetKey: OffsetKey => parseOffsets(offsetKey, record.value) case groupMetadataKey: GroupMetadataKey => parseGroupMetadata(groupMetadataKey, record.value) - case _: UnknownKey => (Some(""), Some("")) + case unknownKey: UnknownKey => (Some(s"UNKNOWN(version=${unknownKey.version})"), None) case _ => throw new KafkaException("Failed to decode message using offset topic decoder (message had an invalid key)") } } @@ -1369,17 +1370,14 @@ trait BaseKey{ } case class OffsetKey(version: Short, key: GroupTopicPartition) extends BaseKey { - override def toString: String = key.toString } case class GroupMetadataKey(version: Short, key: String) extends BaseKey { - override def toString: String = key } -case class UnknownKey(version: Short, key: String = null) extends BaseKey { - +case class UnknownKey(version: Short) extends BaseKey { + override def key: String = null override def toString: String = key } - diff --git a/core/src/main/scala/kafka/coordinator/transaction/TransactionLog.scala b/core/src/main/scala/kafka/coordinator/transaction/TransactionLog.scala index 4623bb8a1a37b..7542a77ed69f3 100644 --- a/core/src/main/scala/kafka/coordinator/transaction/TransactionLog.scala +++ b/core/src/main/scala/kafka/coordinator/transaction/TransactionLog.scala @@ -20,7 +20,6 @@ import java.io.PrintStream import java.nio.ByteBuffer import java.nio.charset.StandardCharsets import kafka.internals.generated.{TransactionLogKey, TransactionLogValue} -import kafka.utils.Logging import org.apache.kafka.clients.consumer.ConsumerRecord import org.apache.kafka.common.protocol.{ByteBufferAccessor, MessageUtil} import org.apache.kafka.common.record.{CompressionType, Record, RecordBatch} @@ -37,7 +36,7 @@ import scala.jdk.CollectionConverters._ * key version 0: [transactionalId] * -> value version 0: [producer_id, producer_epoch, expire_timestamp, status, [topic, [partition] ], timestamp] */ -object TransactionLog extends Logging { +object TransactionLog { // log-level config default values and enforced values val DefaultNumPartitions: Int = 50 @@ -107,8 +106,6 @@ object TransactionLog extends Logging { transactionalId = value.transactionalId ) } else { - warn(s"Unknown version $version from the transaction log message." + - s" The downgraded coordinator will ignore this key and corresponding value.") UnknownKey(version) } } @@ -213,7 +210,8 @@ case class TxnKey(version: Short, transactionalId: String) extends BaseKey { override def toString: String = transactionalId } -case class UnknownKey(version: Short, transactionalId: String = null) extends BaseKey { +case class UnknownKey(version: Short) extends BaseKey { + override def transactionalId: String = null override def toString: String = transactionalId } diff --git a/core/src/main/scala/kafka/coordinator/transaction/TransactionStateManager.scala b/core/src/main/scala/kafka/coordinator/transaction/TransactionStateManager.scala index 85b0e472b4eda..e836c365ddd0e 100644 --- a/core/src/main/scala/kafka/coordinator/transaction/TransactionStateManager.scala +++ b/core/src/main/scala/kafka/coordinator/transaction/TransactionStateManager.scala @@ -479,7 +479,10 @@ class TransactionStateManager(brokerId: Int, } currOffset = batch.nextOffset - case _: UnknownKey => // ignore unknown keys + case unknownKey: UnknownKey => + warn(s"Unknown message key with version ${unknownKey.version}" + + s" while loading transaction state. Ignoring it. " + + s"It could be a left over from an aborted upgrade.") case unexpectedKey => throw new IllegalStateException(s"Found unexpected key $unexpectedKey while reading transaction log.") diff --git a/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala b/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala index 2329e841e3dce..9c90185980aa2 100644 --- a/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala @@ -642,7 +642,8 @@ class GroupMetadataManagerTest { val groupMetadataRecord = buildStableGroupRecordWithMember(generation, protocolType, protocol, memberId) // Should ignore unknown record - val unknownMessage = MessageUtil.messageWithUnknownVersion() + val unknownKey = new org.apache.kafka.coordinator.group.generated.GroupMetadataKey() + val unknownMessage = MessageUtil.toVersionPrefixedBytes(Short.MaxValue, unknownKey) val unknownRecord = new SimpleRecord(unknownMessage, unknownMessage) val records = MemoryRecords.withRecords(startOffset, CompressionType.NONE, @@ -2772,7 +2773,8 @@ class GroupMetadataManagerTest { def testIgnoreUnknownMessageKeyVersion(): Unit = { val record = new org.apache.kafka.coordinator.group.generated.GroupMetadataKey() val unknownRecord = MessageUtil.toVersionPrefixedBytes(Short.MaxValue, record) - GroupMetadataManager.readMessageKey(ByteBuffer.wrap(unknownRecord)) + val key = GroupMetadataManager.readMessageKey(ByteBuffer.wrap(unknownRecord)) + assertEquals(UnknownKey(Short.MaxValue), key) } } diff --git a/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionLogTest.scala b/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionLogTest.scala index d8e72fe838403..71d440e517318 100644 --- a/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionLogTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionLogTest.scala @@ -142,7 +142,8 @@ class TransactionLogTest { def testReadUnknownMessageKeyVersion(): Unit = { val record = new TransactionLogKey() val unknownRecord = MessageUtil.toVersionPrefixedBytes(Short.MaxValue, record) - TransactionLog.readTxnRecordKey(ByteBuffer.wrap(unknownRecord)) + val key = TransactionLog.readTxnRecordKey(ByteBuffer.wrap(unknownRecord)) + assertEquals(UnknownKey(Short.MaxValue), key) } } diff --git a/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionStateManagerTest.scala b/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionStateManagerTest.scala index dcf31d5c48662..188bc57663883 100644 --- a/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionStateManagerTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionStateManagerTest.scala @@ -16,6 +16,8 @@ */ package kafka.coordinator.transaction +import kafka.internals.generated.TransactionLogKey + import java.lang.management.ManagementFactory import java.nio.ByteBuffer import java.util.concurrent.CountDownLatch @@ -758,7 +760,7 @@ class TransactionStateManagerTest { appendedRecords.values.foreach { batches => batches.foreach { records => records.records.forEach { record => - val transactionalId = TransactionLog.readTxnRecordKey(record.key).get.transactionalId + val transactionalId = TransactionLog.readTxnRecordKey(record.key).transactionalId assertNull(record.value) expiredTransactionalIds += transactionalId assertEquals(Right(None), transactionManager.getTransactionState(transactionalId)) @@ -1096,7 +1098,8 @@ class TransactionStateManagerTest { txnRecords += new SimpleRecord(txnMessageKeyBytes1, TransactionLog.valueToBytes(txnMetadata1.prepareNoTransit())) val startOffset = 0L - val unknownMessage = MessageUtil.messageWithUnknownVersion() + val unknownKey = new TransactionLogKey() + val unknownMessage = MessageUtil.toVersionPrefixedBytes(Short.MaxValue, unknownKey) val unknownRecord = new SimpleRecord(unknownMessage, unknownMessage) val records = MemoryRecords.withRecords(startOffset, CompressionType.NONE, From a3fcd3685133d06af89698b4b4c7edf18bbabbec Mon Sep 17 00:00:00 2001 From: Jeff Kim Date: Fri, 14 Apr 2023 12:29:46 -0400 Subject: [PATCH 4/6] address comments --- .../coordinator/group/GroupMetadataManager.scala | 13 +++---------- .../coordinator/transaction/TransactionLog.scala | 15 +++++---------- .../transaction/TransactionStateManager.scala | 3 --- .../transaction/TransactionStateManagerTest.scala | 1 - 4 files changed, 8 insertions(+), 24 deletions(-) diff --git a/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala b/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala index 4ce684aa6b1ca..00f555b82689d 100644 --- a/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala +++ b/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala @@ -576,10 +576,7 @@ class GroupMetadataManager(brokerId: Int, } } - // Visible for testing - private[group] def doLoadGroupsAndOffsets(topicPartition: TopicPartition, - onGroupLoaded: GroupMetadata => Unit): Unit = { - + private def doLoadGroupsAndOffsets(topicPartition: TopicPartition, onGroupLoaded: GroupMetadata => Unit): Unit = { def logEndOffset: Long = replicaManager.getLogEndOffset(topicPartition).getOrElse(-1L) replicaManager.getLog(topicPartition) match { @@ -690,9 +687,6 @@ class GroupMetadataManager(brokerId: Int, warn(s"Unknown message key with version ${unknownKey.version}" + s" while loading offsets and group metadata. Ignoring it. " + s"It could be a left over from an aborted upgrade.") - - case unexpectedKey => - throw new IllegalStateException(s"Unexpected message key $unexpectedKey while loading offsets and group metadata") } } } @@ -1285,8 +1279,7 @@ object GroupMetadataManager { GroupMetadataManager.readMessageKey(record.key) match { case offsetKey: OffsetKey => parseOffsets(offsetKey, record.value) case groupMetadataKey: GroupMetadataKey => parseGroupMetadata(groupMetadataKey, record.value) - case unknownKey: UnknownKey => (Some(s"UNKNOWN(version=${unknownKey.version})"), None) - case _ => throw new KafkaException("Failed to decode message using offset topic decoder (message had an invalid key)") + case unknownKey: UnknownKey => (Some(s"unknown::version=${unknownKey.version}"), None) } } } @@ -1364,7 +1357,7 @@ case class GroupTopicPartition(group: String, topicPartition: TopicPartition) { "[%s,%s,%d]".format(group, topicPartition.topic, topicPartition.partition) } -trait BaseKey{ +sealed trait BaseKey{ def version: Short def key: Any } diff --git a/core/src/main/scala/kafka/coordinator/transaction/TransactionLog.scala b/core/src/main/scala/kafka/coordinator/transaction/TransactionLog.scala index 7542a77ed69f3..30bd517c0093e 100644 --- a/core/src/main/scala/kafka/coordinator/transaction/TransactionLog.scala +++ b/core/src/main/scala/kafka/coordinator/transaction/TransactionLog.scala @@ -162,10 +162,8 @@ object TransactionLog { output.write(producerIdMetadata.getOrElse("NULL").toString.getBytes(StandardCharsets.UTF_8)) output.write("\n".getBytes(StandardCharsets.UTF_8)) - case _: UnknownKey => // Only print if this message is a transaction record - - case unexpectedKey => - throw new IllegalStateException(s"Found unexpected key $unexpectedKey while reading transaction log.") + case unknownKey: UnknownKey => + output.write(s"unknown::version=${unknownKey.version}\n".getBytes(StandardCharsets.UTF_8)) } } } @@ -191,17 +189,14 @@ object TransactionLog { (Some(keyString), Some(valueString)) - case _: UnknownKey => - (Some(""), Some("")) - - case unexpectedKey => - throw new IllegalStateException(s"Found unexpected key $unexpectedKey while formatting transaction log.") + case unknownKey: UnknownKey => + (Some(s"unknown::version=${unknownKey.version}"), None) } } } -trait BaseKey{ +sealed trait BaseKey{ def version: Short def transactionalId: String } diff --git a/core/src/main/scala/kafka/coordinator/transaction/TransactionStateManager.scala b/core/src/main/scala/kafka/coordinator/transaction/TransactionStateManager.scala index e836c365ddd0e..93abc09fb24bd 100644 --- a/core/src/main/scala/kafka/coordinator/transaction/TransactionStateManager.scala +++ b/core/src/main/scala/kafka/coordinator/transaction/TransactionStateManager.scala @@ -483,9 +483,6 @@ class TransactionStateManager(brokerId: Int, warn(s"Unknown message key with version ${unknownKey.version}" + s" while loading transaction state. Ignoring it. " + s"It could be a left over from an aborted upgrade.") - - case unexpectedKey => - throw new IllegalStateException(s"Found unexpected key $unexpectedKey while reading transaction log.") } } } diff --git a/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionStateManagerTest.scala b/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionStateManagerTest.scala index 188bc57663883..0e03c1266aa15 100644 --- a/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionStateManagerTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionStateManagerTest.scala @@ -1088,7 +1088,6 @@ class TransactionStateManagerTest { assertTrue(partitionLoadTime( "partition-load-time-avg") >= 0) } - @Test def testIgnoreUnknownRecordType(): Unit = { txnMetadata1.state = PrepareCommit From 92718230bc18a9078c6ad5690d6b217b36245824 Mon Sep 17 00:00:00 2001 From: Jeff Kim Date: Fri, 14 Apr 2023 16:43:44 -0400 Subject: [PATCH 5/6] tidy --- .../kafka/coordinator/group/GroupMetadataManager.scala | 6 +++--- .../kafka/coordinator/group/GroupMetadataManagerTest.scala | 1 - .../kafka/coordinator/transaction/TransactionLogTest.scala | 1 - 3 files changed, 3 insertions(+), 5 deletions(-) diff --git a/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala b/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala index 00f555b82689d..acc154fe4cc67 100644 --- a/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala +++ b/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala @@ -1277,9 +1277,9 @@ object GroupMetadataManager { throw new KafkaException("Failed to decode message using offset topic decoder (message had a missing key)") } else { GroupMetadataManager.readMessageKey(record.key) match { - case offsetKey: OffsetKey => parseOffsets(offsetKey, record.value) - case groupMetadataKey: GroupMetadataKey => parseGroupMetadata(groupMetadataKey, record.value) - case unknownKey: UnknownKey => (Some(s"unknown::version=${unknownKey.version}"), None) + case offsetKey: OffsetKey => parseOffsets(offsetKey, record.value) + case groupMetadataKey: GroupMetadataKey => parseGroupMetadata(groupMetadataKey, record.value) + case unknownKey: UnknownKey => (Some(s"unknown::version=${unknownKey.version}"), None) } } } diff --git a/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala b/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala index 9c90185980aa2..0147cf661874e 100644 --- a/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala @@ -2776,5 +2776,4 @@ class GroupMetadataManagerTest { val key = GroupMetadataManager.readMessageKey(ByteBuffer.wrap(unknownRecord)) assertEquals(UnknownKey(Short.MaxValue), key) } - } diff --git a/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionLogTest.scala b/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionLogTest.scala index 71d440e517318..5e1e5988e69da 100644 --- a/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionLogTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionLogTest.scala @@ -145,5 +145,4 @@ class TransactionLogTest { val key = TransactionLog.readTxnRecordKey(ByteBuffer.wrap(unknownRecord)) assertEquals(UnknownKey(Short.MaxValue), key) } - } From 7244a25a091e9ab4d006b698bb9a0f054dd0826c Mon Sep 17 00:00:00 2001 From: Jeff Kim Date: Mon, 17 Apr 2023 10:16:23 -0400 Subject: [PATCH 6/6] address comments --- .../group/GroupMetadataManager.scala | 5 +- .../transaction/TransactionStateManager.scala | 4 +- .../group/GroupMetadataManagerTest.scala | 57 ++++++++++++++++--- .../transaction/TransactionLogTest.scala | 2 +- 4 files changed, 55 insertions(+), 13 deletions(-) diff --git a/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala b/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala index acc154fe4cc67..7af924f9edc94 100644 --- a/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala +++ b/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala @@ -683,10 +683,9 @@ class GroupMetadataManager(brokerId: Int, } case unknownKey: UnknownKey => - // Unknown versions may exist when a downgraded coordinator is reading records from the log. warn(s"Unknown message key with version ${unknownKey.version}" + - s" while loading offsets and group metadata. Ignoring it. " + - s"It could be a left over from an aborted upgrade.") + s" while loading offsets and group metadata from $topicPartition. Ignoring it. " + + "It could be a left over from an aborted upgrade.") } } } diff --git a/core/src/main/scala/kafka/coordinator/transaction/TransactionStateManager.scala b/core/src/main/scala/kafka/coordinator/transaction/TransactionStateManager.scala index 93abc09fb24bd..d460edb029dda 100644 --- a/core/src/main/scala/kafka/coordinator/transaction/TransactionStateManager.scala +++ b/core/src/main/scala/kafka/coordinator/transaction/TransactionStateManager.scala @@ -481,8 +481,8 @@ class TransactionStateManager(brokerId: Int, case unknownKey: UnknownKey => warn(s"Unknown message key with version ${unknownKey.version}" + - s" while loading transaction state. Ignoring it. " + - s"It could be a left over from an aborted upgrade.") + s" while loading transaction state from $topicPartition. Ignoring it. " + + "It could be a left over from an aborted upgrade.") } } } diff --git a/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala b/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala index 0147cf661874e..185a48020cb96 100644 --- a/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala @@ -641,13 +641,8 @@ class GroupMetadataManagerTest { val memberId = "98098230493" val groupMetadataRecord = buildStableGroupRecordWithMember(generation, protocolType, protocol, memberId) - // Should ignore unknown record - val unknownKey = new org.apache.kafka.coordinator.group.generated.GroupMetadataKey() - val unknownMessage = MessageUtil.toVersionPrefixedBytes(Short.MaxValue, unknownKey) - val unknownRecord = new SimpleRecord(unknownMessage, unknownMessage) - val records = MemoryRecords.withRecords(startOffset, CompressionType.NONE, - (offsetCommitRecords ++ Seq(unknownRecord) ++ Seq(groupMetadataRecord)).toArray: _*) + (offsetCommitRecords ++ Seq(groupMetadataRecord)).toArray: _*) expectGroupMetadataLoad(groupMetadataTopicPartition, startOffset, records) @@ -2770,10 +2765,58 @@ class GroupMetadataManagerTest { } @Test - def testIgnoreUnknownMessageKeyVersion(): Unit = { + def testReadMessageKeyCanReadUnknownMessage(): Unit = { val record = new org.apache.kafka.coordinator.group.generated.GroupMetadataKey() val unknownRecord = MessageUtil.toVersionPrefixedBytes(Short.MaxValue, record) val key = GroupMetadataManager.readMessageKey(ByteBuffer.wrap(unknownRecord)) assertEquals(UnknownKey(Short.MaxValue), key) } + + @Test + def testLoadGroupsAndOffsetsWillIgnoreUnknownMessage(): Unit = { + val generation = 935 + val protocolType = "consumer" + val protocol = "range" + val startOffset = 15L + val committedOffsets = Map( + new TopicPartition("foo", 0) -> 23L, + new TopicPartition("foo", 1) -> 455L, + new TopicPartition("bar", 0) -> 8992L + ) + + val offsetCommitRecords = createCommittedOffsetRecords(committedOffsets) + val memberId = "98098230493" + val groupMetadataRecord = buildStableGroupRecordWithMember(generation, protocolType, protocol, memberId) + + // Should ignore unknown record + val unknownKey = new org.apache.kafka.coordinator.group.generated.GroupMetadataKey() + val lowestUnsupportedVersion = (org.apache.kafka.coordinator.group.generated.GroupMetadataKey + .HIGHEST_SUPPORTED_VERSION + 1).toShort + + val unknownMessage1 = MessageUtil.toVersionPrefixedBytes(Short.MaxValue, unknownKey) + val unknownMessage2 = MessageUtil.toVersionPrefixedBytes(lowestUnsupportedVersion, unknownKey) + val unknownRecord1 = new SimpleRecord(unknownMessage1, unknownMessage1) + val unknownRecord2 = new SimpleRecord(unknownMessage2, unknownMessage2) + + val records = MemoryRecords.withRecords(startOffset, CompressionType.NONE, + (offsetCommitRecords ++ Seq(unknownRecord1, unknownRecord2) ++ Seq(groupMetadataRecord)).toArray: _*) + + expectGroupMetadataLoad(groupTopicPartition, startOffset, records) + + groupMetadataManager.loadGroupsAndOffsets(groupTopicPartition, 1, _ => (), 0L) + + val group = groupMetadataManager.getGroup(groupId).getOrElse(throw new AssertionError("Group was not loaded into the cache")) + assertEquals(groupId, group.groupId) + assertEquals(Stable, group.currentState) + assertEquals(memberId, group.leaderOrNull) + assertEquals(generation, group.generationId) + assertEquals(Some(protocolType), group.protocolType) + assertEquals(protocol, group.protocolName.orNull) + assertEquals(Set(memberId), group.allMembers) + assertEquals(committedOffsets.size, group.allOffsets.size) + committedOffsets.foreach { case (topicPartition, offset) => + assertEquals(Some(offset), group.offset(topicPartition).map(_.offset)) + assertTrue(group.offset(topicPartition).map(_.expireTimestamp).contains(None)) + } + } } diff --git a/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionLogTest.scala b/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionLogTest.scala index 5e1e5988e69da..eb1284278d29b 100644 --- a/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionLogTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionLogTest.scala @@ -139,7 +139,7 @@ class TransactionLogTest { } @Test - def testReadUnknownMessageKeyVersion(): Unit = { + def testReadTxnRecordKeyCanReadUnknownMessage(): Unit = { val record = new TransactionLogKey() val unknownRecord = MessageUtil.toVersionPrefixedBytes(Short.MaxValue, record) val key = TransactionLog.readTxnRecordKey(ByteBuffer.wrap(unknownRecord))