diff --git a/build.gradle b/build.gradle index 012bb4fb0ac1b..5e15200f1474c 100644 --- a/build.gradle +++ b/build.gradle @@ -1827,7 +1827,8 @@ project(':share-coordinator') { args = [ "-p", "org.apache.kafka.coordinator.share.generated", "-o", "${projectDir}/build/generated/main/java/org/apache/kafka/coordinator/share/generated", "-i", "src/main/resources/common/message", - "-m", "MessageDataGenerator", "JsonConverterGenerator" + "-m", "MessageDataGenerator", "JsonConverterGenerator", + "-t", "CoordinatorRecordTypeGenerator" ] inputs.dir("src/main/resources/common/message") .withPropertyName("messages") diff --git a/core/src/test/scala/unit/kafka/tools/DumpLogSegmentsTest.scala b/core/src/test/scala/unit/kafka/tools/DumpLogSegmentsTest.scala index 8aaeb89246622..7fad4216c0e06 100644 --- a/core/src/test/scala/unit/kafka/tools/DumpLogSegmentsTest.scala +++ b/core/src/test/scala/unit/kafka/tools/DumpLogSegmentsTest.scala @@ -41,8 +41,8 @@ import org.apache.kafka.common.utils.{Exit, Utils} import org.apache.kafka.coordinator.common.runtime.CoordinatorRecord import org.apache.kafka.coordinator.group.GroupCoordinatorRecordSerde import org.apache.kafka.coordinator.group.generated.{ConsumerGroupMemberMetadataValue, ConsumerGroupMetadataKey, ConsumerGroupMetadataValue, GroupMetadataKey, GroupMetadataValue} -import org.apache.kafka.coordinator.share.generated.{ShareSnapshotKey, ShareSnapshotValue, ShareUpdateKey, ShareUpdateValue} -import org.apache.kafka.coordinator.share.{ShareCoordinator, ShareCoordinatorRecordSerde} +import org.apache.kafka.coordinator.share.generated.{CoordinatorRecordType, ShareSnapshotKey, ShareSnapshotValue, ShareUpdateKey, ShareUpdateValue} +import org.apache.kafka.coordinator.share.ShareCoordinatorRecordSerde import org.apache.kafka.coordinator.transaction.generated.{TransactionLogKey, TransactionLogValue} import org.apache.kafka.coordinator.transaction.{TransactionCoordinatorRecordSerde, TransactionLogConfig} import org.apache.kafka.metadata.MetadataRecordSerde @@ -1119,7 +1119,7 @@ class DumpLogSegmentsTest { .setGroupId("gs1") .setTopicId(Uuid.fromString("Uj5wn_FqTXirEASvVZRY1w")) .setPartition(0), - ShareCoordinator.SHARE_SNAPSHOT_RECORD_KEY_VERSION), + CoordinatorRecordType.SHARE_SNAPSHOT.id()), new ApiMessageAndVersion(new ShareSnapshotValue() .setSnapshotEpoch(0) .setStateEpoch(0) @@ -1132,7 +1132,7 @@ class DumpLogSegmentsTest { .setDeliveryState(2) .setDeliveryCount(1) ).asJava), - ShareCoordinator.SHARE_SNAPSHOT_RECORD_VALUE_VERSION) + 0.toShort) )) ) @@ -1147,7 +1147,7 @@ class DumpLogSegmentsTest { .setGroupId("gs1") .setTopicId(Uuid.fromString("Uj5wn_FqTXirEASvVZRY1w")) .setPartition(0), - ShareCoordinator.SHARE_UPDATE_RECORD_KEY_VERSION), + CoordinatorRecordType.SHARE_UPDATE.id()), new ApiMessageAndVersion(new ShareUpdateValue() .setSnapshotEpoch(0) .setLeaderEpoch(0) @@ -1175,7 +1175,7 @@ class DumpLogSegmentsTest { .setGroupId("gs1") .setTopicId(Uuid.fromString("Uj5wn_FqTXirEASvVZRY1w")) .setPartition(0), - 0.toShort + CoordinatorRecordType.SHARE_SNAPSHOT.id() ), null )) diff --git a/share-coordinator/src/main/java/org/apache/kafka/coordinator/share/ShareCoordinator.java b/share-coordinator/src/main/java/org/apache/kafka/coordinator/share/ShareCoordinator.java index dd56503dbafe5..bbbac426fe99a 100644 --- a/share-coordinator/src/main/java/org/apache/kafka/coordinator/share/ShareCoordinator.java +++ b/share-coordinator/src/main/java/org/apache/kafka/coordinator/share/ShareCoordinator.java @@ -32,11 +32,6 @@ import java.util.function.IntSupplier; public interface ShareCoordinator { - short SHARE_SNAPSHOT_RECORD_KEY_VERSION = 0; - short SHARE_SNAPSHOT_RECORD_VALUE_VERSION = 0; - short SHARE_UPDATE_RECORD_KEY_VERSION = 1; - short SHARE_UPDATE_RECORD_VALUE_VERSION = 1; - /** * Return the partition index for the given key. * diff --git a/share-coordinator/src/main/java/org/apache/kafka/coordinator/share/ShareCoordinatorRecordHelpers.java b/share-coordinator/src/main/java/org/apache/kafka/coordinator/share/ShareCoordinatorRecordHelpers.java index bd4bd57a34a34..abde3f442b008 100644 --- a/share-coordinator/src/main/java/org/apache/kafka/coordinator/share/ShareCoordinatorRecordHelpers.java +++ b/share-coordinator/src/main/java/org/apache/kafka/coordinator/share/ShareCoordinatorRecordHelpers.java @@ -18,6 +18,7 @@ import org.apache.kafka.common.Uuid; import org.apache.kafka.coordinator.common.runtime.CoordinatorRecord; +import org.apache.kafka.coordinator.share.generated.CoordinatorRecordType; import org.apache.kafka.coordinator.share.generated.ShareSnapshotKey; import org.apache.kafka.coordinator.share.generated.ShareSnapshotValue; import org.apache.kafka.coordinator.share.generated.ShareUpdateKey; @@ -33,7 +34,7 @@ public static CoordinatorRecord newShareSnapshotRecord(String groupId, Uuid topi .setGroupId(groupId) .setTopicId(topicId) .setPartition(partitionId), - ShareCoordinator.SHARE_SNAPSHOT_RECORD_KEY_VERSION), + CoordinatorRecordType.SHARE_SNAPSHOT.id()), new ApiMessageAndVersion(new ShareSnapshotValue() .setSnapshotEpoch(offsetData.snapshotEpoch()) .setStateEpoch(offsetData.stateEpoch()) @@ -46,7 +47,7 @@ public static CoordinatorRecord newShareSnapshotRecord(String groupId, Uuid topi .setDeliveryCount(batch.deliveryCount()) .setDeliveryState(batch.deliveryState())) .collect(Collectors.toList())), - ShareCoordinator.SHARE_SNAPSHOT_RECORD_VALUE_VERSION) + (short) 0) ); } @@ -56,7 +57,7 @@ public static CoordinatorRecord newShareSnapshotUpdateRecord(String groupId, Uui .setGroupId(groupId) .setTopicId(topicId) .setPartition(partitionId), - ShareCoordinator.SHARE_UPDATE_RECORD_KEY_VERSION), + CoordinatorRecordType.SHARE_UPDATE.id()), new ApiMessageAndVersion(new ShareUpdateValue() .setSnapshotEpoch(offsetData.snapshotEpoch()) .setLeaderEpoch(offsetData.leaderEpoch()) @@ -68,7 +69,7 @@ public static CoordinatorRecord newShareSnapshotUpdateRecord(String groupId, Uui .setDeliveryCount(batch.deliveryCount()) .setDeliveryState(batch.deliveryState())) .collect(Collectors.toList())), - ShareCoordinator.SHARE_UPDATE_RECORD_VALUE_VERSION) + (short) 0) ); } } diff --git a/share-coordinator/src/main/java/org/apache/kafka/coordinator/share/ShareCoordinatorRecordSerde.java b/share-coordinator/src/main/java/org/apache/kafka/coordinator/share/ShareCoordinatorRecordSerde.java index a620289d17ae4..28f59e57d336a 100644 --- a/share-coordinator/src/main/java/org/apache/kafka/coordinator/share/ShareCoordinatorRecordSerde.java +++ b/share-coordinator/src/main/java/org/apache/kafka/coordinator/share/ShareCoordinatorRecordSerde.java @@ -29,9 +29,9 @@ public class ShareCoordinatorRecordSerde extends CoordinatorRecordSerde { @Override protected ApiMessage apiMessageKeyFor(short recordVersion) { switch (recordVersion) { - case ShareCoordinator.SHARE_SNAPSHOT_RECORD_KEY_VERSION: + case 0: return new ShareSnapshotKey(); - case ShareCoordinator.SHARE_UPDATE_RECORD_KEY_VERSION: + case 1: return new ShareUpdateKey(); default: throw new CoordinatorLoader.UnknownRecordTypeException(recordVersion); @@ -41,9 +41,9 @@ protected ApiMessage apiMessageKeyFor(short recordVersion) { @Override protected ApiMessage apiMessageValueFor(short recordVersion) { switch (recordVersion) { - case ShareCoordinator.SHARE_SNAPSHOT_RECORD_VALUE_VERSION: + case 0: return new ShareSnapshotValue(); - case ShareCoordinator.SHARE_UPDATE_RECORD_VALUE_VERSION: + case 1: return new ShareUpdateValue(); default: throw new CoordinatorLoader.UnknownRecordTypeException(recordVersion); diff --git a/share-coordinator/src/main/java/org/apache/kafka/coordinator/share/ShareCoordinatorShard.java b/share-coordinator/src/main/java/org/apache/kafka/coordinator/share/ShareCoordinatorShard.java index 6f2a9c27b7419..ab867d331aa59 100644 --- a/share-coordinator/src/main/java/org/apache/kafka/coordinator/share/ShareCoordinatorShard.java +++ b/share-coordinator/src/main/java/org/apache/kafka/coordinator/share/ShareCoordinatorShard.java @@ -19,6 +19,7 @@ import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.Uuid; +import org.apache.kafka.common.errors.UnsupportedVersionException; import org.apache.kafka.common.message.ReadShareGroupStateRequestData; import org.apache.kafka.common.message.ReadShareGroupStateResponseData; import org.apache.kafka.common.message.WriteShareGroupStateRequestData; @@ -38,6 +39,7 @@ import org.apache.kafka.coordinator.common.runtime.CoordinatorShard; import org.apache.kafka.coordinator.common.runtime.CoordinatorShardBuilder; import org.apache.kafka.coordinator.common.runtime.CoordinatorTimer; +import org.apache.kafka.coordinator.share.generated.CoordinatorRecordType; import org.apache.kafka.coordinator.share.generated.ShareSnapshotKey; import org.apache.kafka.coordinator.share.generated.ShareSnapshotValue; import org.apache.kafka.coordinator.share.generated.ShareUpdateKey; @@ -206,15 +208,19 @@ public void replay(long offset, long producerId, short producerEpoch, Coordinato ApiMessageAndVersion key = record.key(); ApiMessageAndVersion value = record.value(); - switch (key.version()) { - case ShareCoordinator.SHARE_SNAPSHOT_RECORD_KEY_VERSION: // ShareSnapshot - handleShareSnapshot((ShareSnapshotKey) key.message(), (ShareSnapshotValue) messageOrNull(value), offset); - break; - case ShareCoordinator.SHARE_UPDATE_RECORD_KEY_VERSION: // ShareUpdate - handleShareUpdate((ShareUpdateKey) key.message(), (ShareUpdateValue) messageOrNull(value)); - break; - default: - // Noop + try { + switch (CoordinatorRecordType.fromId(key.version())) { + case SHARE_SNAPSHOT: + handleShareSnapshot((ShareSnapshotKey) key.message(), (ShareSnapshotValue) messageOrNull(value), offset); + break; + case SHARE_UPDATE: + handleShareUpdate((ShareUpdateKey) key.message(), (ShareUpdateValue) messageOrNull(value)); + break; + default: + // Noop + } + } catch (UnsupportedVersionException ex) { + // Ignore } } diff --git a/share-coordinator/src/main/resources/common/message/ShareSnapshotKey.json b/share-coordinator/src/main/resources/common/message/ShareSnapshotKey.json index f8de8a237feed..128f39a0182f8 100644 --- a/share-coordinator/src/main/resources/common/message/ShareSnapshotKey.json +++ b/share-coordinator/src/main/resources/common/message/ShareSnapshotKey.json @@ -16,7 +16,8 @@ // KIP-932 is in development. This schema is subject to non-backwards-compatible changes. { - "type": "data", + "apiKey": 0, + "type": "coordinator-key", "name": "ShareSnapshotKey", "validVersions": "0", "flexibleVersions": "none", diff --git a/share-coordinator/src/main/resources/common/message/ShareSnapshotValue.json b/share-coordinator/src/main/resources/common/message/ShareSnapshotValue.json index 0fceee616b300..4b375381c3b5c 100644 --- a/share-coordinator/src/main/resources/common/message/ShareSnapshotValue.json +++ b/share-coordinator/src/main/resources/common/message/ShareSnapshotValue.json @@ -16,7 +16,8 @@ // KIP-932 is in development. This schema is subject to non-backwards-compatible changes. { - "type": "data", + "apiKey": 0, + "type": "coordinator-value", "name": "ShareSnapshotValue", "validVersions": "0", "flexibleVersions": "0+", diff --git a/share-coordinator/src/main/resources/common/message/ShareUpdateKey.json b/share-coordinator/src/main/resources/common/message/ShareUpdateKey.json index 923ed5cf6f6af..ba919dc6b1c08 100644 --- a/share-coordinator/src/main/resources/common/message/ShareUpdateKey.json +++ b/share-coordinator/src/main/resources/common/message/ShareUpdateKey.json @@ -16,16 +16,17 @@ // KIP-932 is in development. This schema is subject to non-backwards-compatible changes. { - "type": "data", + "apiKey": 1, + "type": "coordinator-key", "name": "ShareUpdateKey", - "validVersions": "1", + "validVersions": "0", "flexibleVersions": "none", "fields": [ - { "name": "GroupId", "type": "string", "versions": "1", + { "name": "GroupId", "type": "string", "versions": "0", "about": "The group id." }, - { "name": "TopicId", "type": "uuid", "versions": "1", + { "name": "TopicId", "type": "uuid", "versions": "0", "about": "The topic id." }, - { "name": "Partition", "type": "int32", "versions": "1", + { "name": "Partition", "type": "int32", "versions": "0", "about": "The partition index." } ] } diff --git a/share-coordinator/src/main/resources/common/message/ShareUpdateValue.json b/share-coordinator/src/main/resources/common/message/ShareUpdateValue.json index 6af282db6fd07..389cf688c8709 100644 --- a/share-coordinator/src/main/resources/common/message/ShareUpdateValue.json +++ b/share-coordinator/src/main/resources/common/message/ShareUpdateValue.json @@ -16,7 +16,8 @@ // KIP-932 is in development. This schema is subject to non-backwards-compatible changes. { - "type": "data", + "apiKey": 1, + "type": "coordinator-value", "name": "ShareUpdateValue", "validVersions": "0", "flexibleVersions": "0+", diff --git a/share-coordinator/src/test/java/org/apache/kafka/coordinator/share/ShareCoordinatorRecordHelpersTest.java b/share-coordinator/src/test/java/org/apache/kafka/coordinator/share/ShareCoordinatorRecordHelpersTest.java index 1f47cf9c166c4..6b33727ac19ab 100644 --- a/share-coordinator/src/test/java/org/apache/kafka/coordinator/share/ShareCoordinatorRecordHelpersTest.java +++ b/share-coordinator/src/test/java/org/apache/kafka/coordinator/share/ShareCoordinatorRecordHelpersTest.java @@ -57,7 +57,7 @@ public void testNewShareSnapshotRecord() { .setGroupId(groupId) .setTopicId(topicId) .setPartition(partitionId), - ShareCoordinator.SHARE_SNAPSHOT_RECORD_KEY_VERSION), + (short) 0), new ApiMessageAndVersion( new ShareSnapshotValue() .setSnapshotEpoch(0) @@ -70,7 +70,7 @@ public void testNewShareSnapshotRecord() { .setLastOffset(10L) .setDeliveryState((byte) 0) .setDeliveryCount((short) 1))), - ShareCoordinator.SHARE_SNAPSHOT_RECORD_VALUE_VERSION)); + (short) 0)); assertEquals(expectedRecord, record); } @@ -100,7 +100,7 @@ public void testNewShareUpdateRecord() { .setGroupId(groupId) .setTopicId(topicId) .setPartition(partitionId), - ShareCoordinator.SHARE_UPDATE_RECORD_KEY_VERSION), + (short) 1), new ApiMessageAndVersion( new ShareUpdateValue() .setSnapshotEpoch(0) @@ -112,7 +112,7 @@ public void testNewShareUpdateRecord() { .setLastOffset(10L) .setDeliveryState((byte) 0) .setDeliveryCount((short) 1))), - ShareCoordinator.SHARE_UPDATE_RECORD_VALUE_VERSION)); + (short) 0)); assertEquals(expectedRecord, record); } diff --git a/share-coordinator/src/test/java/org/apache/kafka/coordinator/share/ShareCoordinatorRecordSerdeTest.java b/share-coordinator/src/test/java/org/apache/kafka/coordinator/share/ShareCoordinatorRecordSerdeTest.java index c62abdb13dc99..c11df0c19bb11 100644 --- a/share-coordinator/src/test/java/org/apache/kafka/coordinator/share/ShareCoordinatorRecordSerdeTest.java +++ b/share-coordinator/src/test/java/org/apache/kafka/coordinator/share/ShareCoordinatorRecordSerdeTest.java @@ -21,6 +21,7 @@ import org.apache.kafka.common.protocol.MessageUtil; import org.apache.kafka.coordinator.common.runtime.CoordinatorLoader; import org.apache.kafka.coordinator.common.runtime.CoordinatorRecord; +import org.apache.kafka.coordinator.share.generated.CoordinatorRecordType; import org.apache.kafka.coordinator.share.generated.ShareSnapshotKey; import org.apache.kafka.coordinator.share.generated.ShareSnapshotValue; import org.apache.kafka.coordinator.share.generated.ShareUpdateKey; @@ -75,7 +76,7 @@ public void testSerializeNullValue() { .setGroupId("group") .setTopicId(Uuid.randomUuid()) .setPartition(1), - ShareCoordinator.SHARE_SNAPSHOT_RECORD_KEY_VERSION + CoordinatorRecordType.SHARE_SNAPSHOT.id() ), null ); @@ -104,7 +105,7 @@ public void testDeserializeWithTombstoneForValue() { .setGroupId("groupId") .setTopicId(Uuid.randomUuid()) .setPartition(1), - ShareCoordinator.SHARE_SNAPSHOT_RECORD_KEY_VERSION + CoordinatorRecordType.SHARE_SNAPSHOT.id() ); ByteBuffer keyBuffer = MessageUtil.toVersionPrefixedByteBuffer(key.version(), key.message()); @@ -145,7 +146,7 @@ public void testDeserializeWithValueEmptyBuffer() { .setGroupId("foo") .setTopicId(Uuid.randomUuid()) .setPartition(1), - ShareCoordinator.SHARE_SNAPSHOT_RECORD_KEY_VERSION + CoordinatorRecordType.SHARE_SNAPSHOT.id() ); ByteBuffer keyBuffer = MessageUtil.toVersionPrefixedByteBuffer(key.version(), key.message()); @@ -181,7 +182,7 @@ public void testDeserializeWithInvalidValueBytes() { .setGroupId("foo") .setTopicId(Uuid.randomUuid()) .setPartition(1), - ShareCoordinator.SHARE_SNAPSHOT_RECORD_KEY_VERSION + CoordinatorRecordType.SHARE_SNAPSHOT.id() ); ByteBuffer keyBuffer = MessageUtil.toVersionPrefixedByteBuffer(key.version(), key.message()); @@ -228,7 +229,7 @@ private static CoordinatorRecord getShareSnapshotRecord(String groupId, Uuid top .setGroupId(groupId) .setTopicId(topicId) .setPartition(partitionId), - ShareCoordinator.SHARE_SNAPSHOT_RECORD_KEY_VERSION + CoordinatorRecordType.SHARE_SNAPSHOT.id() ), new ApiMessageAndVersion( new ShareSnapshotValue() diff --git a/tools/src/main/java/org/apache/kafka/tools/consumer/group/share/ShareGroupStateMessageFormatter.java b/tools/src/main/java/org/apache/kafka/tools/consumer/group/share/ShareGroupStateMessageFormatter.java index c11b12d44f949..c5358b62d1b6c 100644 --- a/tools/src/main/java/org/apache/kafka/tools/consumer/group/share/ShareGroupStateMessageFormatter.java +++ b/tools/src/main/java/org/apache/kafka/tools/consumer/group/share/ShareGroupStateMessageFormatter.java @@ -18,8 +18,10 @@ import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.common.MessageFormatter; +import org.apache.kafka.common.errors.UnsupportedVersionException; import org.apache.kafka.common.protocol.ApiMessage; import org.apache.kafka.common.protocol.ByteBufferAccessor; +import org.apache.kafka.coordinator.share.generated.CoordinatorRecordType; import org.apache.kafka.coordinator.share.generated.ShareSnapshotKey; import org.apache.kafka.coordinator.share.generated.ShareSnapshotKeyJsonConverter; import org.apache.kafka.coordinator.share.generated.ShareSnapshotValue; @@ -101,13 +103,16 @@ private JsonNode readToKeyJson(ByteBuffer byteBuffer, short version) { private Optional readToSnapshotMessageKey(ByteBuffer byteBuffer) { short version = byteBuffer.getShort(); - if (version >= ShareSnapshotKey.LOWEST_SUPPORTED_VERSION - && version <= ShareSnapshotKey.HIGHEST_SUPPORTED_VERSION) { - return Optional.of(new ShareSnapshotKey(new ByteBufferAccessor(byteBuffer), version)); - } else if (version >= ShareUpdateKey.LOWEST_SUPPORTED_VERSION - && version <= ShareUpdateKey.HIGHEST_SUPPORTED_VERSION) { - return Optional.of(new ShareUpdateKey(new ByteBufferAccessor(byteBuffer), version)); - } else { + try { + switch (CoordinatorRecordType.fromId(version)) { + case SHARE_SNAPSHOT: + return Optional.of(new ShareSnapshotKey(new ByteBufferAccessor(byteBuffer), version)); + case SHARE_UPDATE: + return Optional.of(new ShareUpdateKey(new ByteBufferAccessor(byteBuffer), version)); + default: + return Optional.empty(); + } + } catch (UnsupportedVersionException ex) { return Optional.empty(); } } @@ -155,13 +160,17 @@ private Optional readToSnapshotMessageValue(ByteBuffer byteBuffer, s // Check the key version here as that will determine which type // of value record to fetch. Both share update and share snapshot // value records can have the same version. - if (keyVersion >= ShareSnapshotKey.LOWEST_SUPPORTED_VERSION - && keyVersion <= ShareSnapshotKey.HIGHEST_SUPPORTED_VERSION) { - return Optional.of(new ShareSnapshotValue(new ByteBufferAccessor(byteBuffer), version)); - } else if (keyVersion >= ShareUpdateKey.LOWEST_SUPPORTED_VERSION - && keyVersion <= ShareUpdateKey.HIGHEST_SUPPORTED_VERSION) { - return Optional.of(new ShareUpdateValue(new ByteBufferAccessor(byteBuffer), version)); + try { + switch (CoordinatorRecordType.fromId(keyVersion)) { + case SHARE_SNAPSHOT: + return Optional.of(new ShareSnapshotValue(new ByteBufferAccessor(byteBuffer), version)); + case SHARE_UPDATE: + return Optional.of(new ShareUpdateValue(new ByteBufferAccessor(byteBuffer), version)); + default: + return Optional.empty(); + } + } catch (UnsupportedVersionException ex) { + return Optional.empty(); } - return Optional.empty(); } }