-
Notifications
You must be signed in to change notification settings - Fork 15.4k
KAFKA-17011: Fix a bug preventing features from supporting v0 #16421
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from 11 commits
1c1e556
6235178
1dc5f99
c0e3c7c
3e07a3b
ab1b566
234ea4b
04ecd25
890c184
6ad6749
c1ba1f5
85a827a
d6eb186
f4bd3cc
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -35,6 +35,7 @@ | |
|
|
||
| import java.nio.ByteBuffer; | ||
| import java.util.Map; | ||
| import java.util.Objects; | ||
| import java.util.Optional; | ||
| import java.util.Set; | ||
|
|
||
|
|
@@ -49,6 +50,71 @@ public class ApiVersionsResponse extends AbstractResponse { | |
|
|
||
| private final ApiVersionsResponseData data; | ||
|
|
||
| public static class Builder { | ||
| private Errors error = Errors.NONE; | ||
| private int throttleTimeMs = 0; | ||
| private ApiVersionCollection apiVersions = null; | ||
| private Features<SupportedVersionRange> supportedFeatures = null; | ||
| private Map<String, Short> finalizedFeatures = null; | ||
| private long finalizedFeaturesEpoch = 0; | ||
| private boolean zkMigrationEnabled = false; | ||
| private boolean suppressFeatureLevel0 = false; | ||
|
|
||
| public Builder setError(Errors error) { | ||
| this.error = error; | ||
| return this; | ||
| } | ||
|
|
||
| public Builder setThrottleTimeMs(int throttleTimeMs) { | ||
| this.throttleTimeMs = throttleTimeMs; | ||
| return this; | ||
| } | ||
|
|
||
| public Builder setApiVersions(ApiVersionCollection apiVersions) { | ||
| this.apiVersions = apiVersions; | ||
| return this; | ||
| } | ||
|
|
||
| public Builder setSupportedFeatures(Features<SupportedVersionRange> supportedFeatures) { | ||
| this.supportedFeatures = supportedFeatures; | ||
| return this; | ||
| } | ||
|
|
||
| public Builder setFinalizedFeatures(Map<String, Short> finalizedFeatures) { | ||
| this.finalizedFeatures = finalizedFeatures; | ||
| return this; | ||
| } | ||
|
|
||
| public Builder setFinalizedFeaturesEpoch(long finalizedFeaturesEpoch) { | ||
| this.finalizedFeaturesEpoch = finalizedFeaturesEpoch; | ||
| return this; | ||
| } | ||
|
|
||
| public Builder setZkMigrationEnabled(boolean zkMigrationEnabled) { | ||
| this.zkMigrationEnabled = zkMigrationEnabled; | ||
| return this; | ||
| } | ||
|
|
||
| public Builder setSuppressFeatureLevel0(boolean suppressFeatureLevel0) { | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. suppressFeatureLevel0 => alterFeatureLevel0 ? |
||
| this.suppressFeatureLevel0 = suppressFeatureLevel0; | ||
| return this; | ||
| } | ||
|
|
||
| public ApiVersionsResponse build() { | ||
| final ApiVersionsResponseData data = new ApiVersionsResponseData(); | ||
| data.setErrorCode(error.code()); | ||
| data.setApiKeys(Objects.requireNonNull(apiVersions)); | ||
| data.setThrottleTimeMs(throttleTimeMs); | ||
| data.setSupportedFeatures( | ||
| maybeFilterSupportedFeatureKeys(Objects.requireNonNull(supportedFeatures), suppressFeatureLevel0)); | ||
| data.setFinalizedFeatures( | ||
| createFinalizedFeatureKeys(Objects.requireNonNull(finalizedFeatures))); | ||
| data.setFinalizedFeaturesEpoch(finalizedFeaturesEpoch); | ||
| data.setZkMigrationReady(zkMigrationEnabled); | ||
| return new ApiVersionsResponse(data); | ||
| } | ||
| } | ||
|
|
||
| public ApiVersionsResponse(ApiVersionsResponseData data) { | ||
| super(ApiKeys.API_VERSIONS); | ||
| this.data = data; | ||
|
|
@@ -105,65 +171,32 @@ public static ApiVersionsResponse parse(ByteBuffer buffer, short version) { | |
| } | ||
| } | ||
|
|
||
| public static ApiVersionsResponse createApiVersionsResponse( | ||
| int throttleTimeMs, | ||
| public static ApiVersionCollection controllerApiVersions( | ||
| RecordVersion minRecordVersion, | ||
| Features<SupportedVersionRange> latestSupportedFeatures, | ||
| Map<String, Short> finalizedFeatures, | ||
| long finalizedFeaturesEpoch, | ||
| NodeApiVersions controllerApiVersions, | ||
| ListenerType listenerType, | ||
| boolean enableUnstableLastVersion, | ||
| boolean zkMigrationEnabled, | ||
| boolean clientTelemetryEnabled | ||
| ) { | ||
| ApiVersionCollection apiKeys; | ||
| if (controllerApiVersions != null) { | ||
| apiKeys = intersectForwardableApis( | ||
| listenerType, | ||
| minRecordVersion, | ||
| controllerApiVersions.allSupportedApiVersions(), | ||
| enableUnstableLastVersion, | ||
| clientTelemetryEnabled | ||
| ); | ||
| } else { | ||
| apiKeys = filterApis( | ||
| minRecordVersion, | ||
| listenerType, | ||
| enableUnstableLastVersion, | ||
| clientTelemetryEnabled | ||
| ); | ||
| } | ||
|
|
||
| return createApiVersionsResponse( | ||
| throttleTimeMs, | ||
| apiKeys, | ||
| latestSupportedFeatures, | ||
| finalizedFeatures, | ||
| finalizedFeaturesEpoch, | ||
| zkMigrationEnabled | ||
| ); | ||
| return intersectForwardableApis( | ||
| listenerType, | ||
| minRecordVersion, | ||
| controllerApiVersions.allSupportedApiVersions(), | ||
| enableUnstableLastVersion, | ||
| clientTelemetryEnabled); | ||
| } | ||
|
|
||
| public static ApiVersionsResponse createApiVersionsResponse( | ||
| int throttleTimeMs, | ||
| ApiVersionCollection apiVersions, | ||
| Features<SupportedVersionRange> latestSupportedFeatures, | ||
| Map<String, Short> finalizedFeatures, | ||
| long finalizedFeaturesEpoch, | ||
| boolean zkMigrationEnabled | ||
| public static ApiVersionCollection brokerApiVersions( | ||
| RecordVersion minRecordVersion, | ||
| ListenerType listenerType, | ||
| boolean enableUnstableLastVersion, | ||
| boolean clientTelemetryEnabled | ||
| ) { | ||
| return new ApiVersionsResponse( | ||
| createApiVersionsResponseData( | ||
| throttleTimeMs, | ||
| Errors.NONE, | ||
| apiVersions, | ||
| latestSupportedFeatures, | ||
| finalizedFeatures, | ||
| finalizedFeaturesEpoch, | ||
| zkMigrationEnabled | ||
| ) | ||
| ); | ||
| return filterApis( | ||
| minRecordVersion, | ||
| listenerType, | ||
| enableUnstableLastVersion, | ||
| clientTelemetryEnabled); | ||
| } | ||
|
|
||
| public static ApiVersionCollection filterApis( | ||
|
|
@@ -249,35 +282,24 @@ public static ApiVersionCollection intersectForwardableApis( | |
| return apiKeys; | ||
| } | ||
|
|
||
| private static ApiVersionsResponseData createApiVersionsResponseData( | ||
| final int throttleTimeMs, | ||
| final Errors error, | ||
| final ApiVersionCollection apiKeys, | ||
| final Features<SupportedVersionRange> latestSupportedFeatures, | ||
| final Map<String, Short> finalizedFeatures, | ||
| final long finalizedFeaturesEpoch, | ||
| final boolean zkMigrationEnabled | ||
| private static SupportedFeatureKeyCollection maybeFilterSupportedFeatureKeys( | ||
| Features<SupportedVersionRange> latestSupportedFeatures, | ||
|
cmccabe marked this conversation as resolved.
|
||
| boolean suppressV0 | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. suppressV0 => alterV0 ? |
||
| ) { | ||
| final ApiVersionsResponseData data = new ApiVersionsResponseData(); | ||
| data.setThrottleTimeMs(throttleTimeMs); | ||
| data.setErrorCode(error.code()); | ||
| data.setApiKeys(apiKeys); | ||
| data.setSupportedFeatures(createSupportedFeatureKeys(latestSupportedFeatures)); | ||
| data.setFinalizedFeatures(createFinalizedFeatureKeys(finalizedFeatures)); | ||
| data.setFinalizedFeaturesEpoch(finalizedFeaturesEpoch); | ||
| data.setZkMigrationReady(zkMigrationEnabled); | ||
|
|
||
| return data; | ||
| } | ||
|
|
||
| private static SupportedFeatureKeyCollection createSupportedFeatureKeys( | ||
| Features<SupportedVersionRange> latestSupportedFeatures) { | ||
| SupportedFeatureKeyCollection converted = new SupportedFeatureKeyCollection(); | ||
| for (Map.Entry<String, SupportedVersionRange> feature : latestSupportedFeatures.features().entrySet()) { | ||
| final SupportedFeatureKey key = new SupportedFeatureKey(); | ||
| final SupportedVersionRange versionRange = feature.getValue(); | ||
| final SupportedFeatureKey key = new SupportedFeatureKey(); | ||
| key.setName(feature.getKey()); | ||
| key.setMinVersion(versionRange.min()); | ||
| if (suppressV0 && versionRange.min() == 0) { | ||
| // Some older clients will have deserialization problems if a feature's | ||
| // minimum supported level is 0. Therefore, when preparing ApiVersionResponse | ||
| // at versions less than 4, we must set the minimum version for these features | ||
| // to 1 rather than 0. See KAFKA-17011 for details. | ||
| key.setMinVersion((short) 1); | ||
| } else { | ||
| key.setMinVersion(versionRange.min()); | ||
| } | ||
| key.setMaxVersion(versionRange.max()); | ||
| converted.add(key); | ||
| } | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -23,6 +23,7 @@ | |
| import org.apache.kafka.common.protocol.Errors; | ||
|
|
||
| import java.nio.ByteBuffer; | ||
| import java.util.Iterator; | ||
|
|
||
| public class BrokerRegistrationRequest extends AbstractRequest { | ||
|
|
||
|
|
@@ -45,7 +46,21 @@ public short oldestAllowedVersion() { | |
|
|
||
| @Override | ||
| public BrokerRegistrationRequest build(short version) { | ||
| return new BrokerRegistrationRequest(data, version); | ||
| if (version < 4) { | ||
| // Workaround for KAFKA-17011: for BrokerRegistrationRequest versions older than 4, | ||
| // translate minSupportedVersion = 0 to minSupportedVersion = 1. | ||
| BrokerRegistrationRequestData newData = data.duplicate(); | ||
| for (Iterator<BrokerRegistrationRequestData.Feature> iter = newData.features().iterator(); | ||
|
chia7712 marked this conversation as resolved.
|
||
| iter.hasNext(); ) { | ||
| BrokerRegistrationRequestData.Feature feature = iter.next(); | ||
| if (feature.minSupportedVersion() == 0) { | ||
| feature.setMinSupportedVersion((short) 1); | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. There is another issue caused by changing the version from 0 to 1. #15685 add the check which assumes the "default version" is 0 Hence, the features having min=1 can NOT pass the check. see following error message:
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. @chia7712 Thanks. Is there a way to reproduce this (e.g. an existing test)?
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
run the following command with #17084
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Originally I thought that we could not set version 0 with the storage or update tools since it would treat it as disabling the feature. But @dajac pointed out to me that this is happening by setting the MV and that picks default features. However, the min version 1 should only be the case for older request versions. I thought we fixed this to so we wouldn't have the requirement for version 4 and for 4.0 where group version is introduced.
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. When I last talked to Colin about this, I believe he said we could not disable the feature if the min version is > 0.
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
@cmccabe : Could you confirm if this is case? How would people upgrade from an old release where a feature is not turned on?
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. I would suspect you have to upgrade the feature before the upgrade to the new release, but we can let Colin confirm.
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Checked with Colin offline. To summarize, if we do increase minVersion in the future, we need to bump up the default version for that feature. An old release may need to first upgrade to a bridge release where the new default version could be set. Given that, we don't need to change the logic in ClusterControlManager.
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
I will sync it to KAFKA-17429 and #17128 |
||
| } | ||
| } | ||
| return new BrokerRegistrationRequest(newData, version); | ||
| } else { | ||
| return new BrokerRegistrationRequest(data, version); | ||
| } | ||
| } | ||
|
|
||
| @Override | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
suppressFeatureLevel0 => alterFeatureLevel0 ?