-
Notifications
You must be signed in to change notification settings - Fork 15.4k
KAFKA-2945: CreateTopic - protocol and server side implementation #1489
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
Closed
Closed
Changes from 1 commit
Commits
Show all changes
5 commits
Select commit
Hold shift + click to select a range
4e1f59a
KAFKA-2945: CreateTopic - protocol and server side implementation
granthenke 390f215
Address Jun's Reviews
granthenke 5ce746f
Address reviews
granthenke 21bf32c
Remove controller.isActive check
granthenke 8bedd81
Add timeout test comment
granthenke File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
31 changes: 31 additions & 0 deletions
31
clients/src/main/java/org/apache/kafka/common/errors/InvalidReplicaAssignmentException.java
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,31 @@ | ||
| /** | ||
| * Licensed to the Apache Software Foundation (ASF) under one or more | ||
| * contributor license agreements. See the NOTICE file distributed with | ||
| * this work for additional information regarding copyright ownership. | ||
| * The ASF licenses this file to You under the Apache License, Version 2.0 | ||
| * (the "License"); you may not use this file except in compliance with | ||
| * the License. You may obtain a copy of the License at | ||
| * | ||
| * http://www.apache.org/licenses/LICENSE-2.0 | ||
| * | ||
| * Unless required by applicable law or agreed to in writing, software | ||
| * distributed under the License is distributed on an "AS IS" BASIS, | ||
| * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. | ||
| * See the License for the specific language governing permissions and | ||
| * limitations under the License. | ||
| */ | ||
| package org.apache.kafka.common.errors; | ||
|
|
||
| public class InvalidReplicaAssignmentException extends ApiException { | ||
|
|
||
| private static final long serialVersionUID = 1L; | ||
|
|
||
| public InvalidReplicaAssignmentException(String message) { | ||
| super(message); | ||
| } | ||
|
|
||
| public InvalidReplicaAssignmentException(String message, Throwable cause) { | ||
| super(message, cause); | ||
| } | ||
|
|
||
| } |
31 changes: 31 additions & 0 deletions
31
clients/src/main/java/org/apache/kafka/common/errors/InvalidReplicationFactorException.java
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,31 @@ | ||
| /** | ||
| * Licensed to the Apache Software Foundation (ASF) under one or more | ||
| * contributor license agreements. See the NOTICE file distributed with | ||
| * this work for additional information regarding copyright ownership. | ||
| * The ASF licenses this file to You under the Apache License, Version 2.0 | ||
| * (the "License"); you may not use this file except in compliance with | ||
| * the License. You may obtain a copy of the License at | ||
| * | ||
| * http://www.apache.org/licenses/LICENSE-2.0 | ||
| * | ||
| * Unless required by applicable law or agreed to in writing, software | ||
| * distributed under the License is distributed on an "AS IS" BASIS, | ||
| * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. | ||
| * See the License for the specific language governing permissions and | ||
| * limitations under the License. | ||
| */ | ||
| package org.apache.kafka.common.errors; | ||
|
|
||
| public class InvalidReplicationFactorException extends ApiException { | ||
|
|
||
| private static final long serialVersionUID = 1L; | ||
|
|
||
| public InvalidReplicationFactorException(String message) { | ||
| super(message); | ||
| } | ||
|
|
||
| public InvalidReplicationFactorException(String message, Throwable cause) { | ||
| super(message, cause); | ||
| } | ||
|
|
||
| } |
36 changes: 36 additions & 0 deletions
36
clients/src/main/java/org/apache/kafka/common/errors/InvalidRequestException.java
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,36 @@ | ||
| /** | ||
| * Licensed to the Apache Software Foundation (ASF) under one or more | ||
| * contributor license agreements. See the NOTICE file distributed with | ||
| * this work for additional information regarding copyright ownership. | ||
| * The ASF licenses this file to You under the Apache License, Version 2.0 | ||
| * (the "License"); you may not use this file except in compliance with | ||
| * the License. You may obtain a copy of the License at | ||
| * | ||
| * http://www.apache.org/licenses/LICENSE-2.0 | ||
| * | ||
| * Unless required by applicable law or agreed to in writing, software | ||
| * distributed under the License is distributed on an "AS IS" BASIS, | ||
| * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. | ||
| * See the License for the specific language governing permissions and | ||
| * limitations under the License. | ||
| */ | ||
| package org.apache.kafka.common.errors; | ||
|
|
||
| /** | ||
| * Thrown when a request breaks basic wire protocol rules. | ||
| * This most likely occurs because of a request being malformed by the client library or | ||
| * the message was sent to an incompatible broker. | ||
| */ | ||
| public class InvalidRequestException extends ApiException { | ||
|
|
||
| private static final long serialVersionUID = 1L; | ||
|
|
||
| public InvalidRequestException(String message) { | ||
| super(message); | ||
| } | ||
|
|
||
| public InvalidRequestException(String message, Throwable cause) { | ||
| super(message, cause); | ||
| } | ||
|
|
||
| } |
28 changes: 28 additions & 0 deletions
28
clients/src/main/java/org/apache/kafka/common/errors/NotControllerException.java
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,28 @@ | ||
| /** | ||
| * Licensed to the Apache Software Foundation (ASF) under one or more contributor license agreements. See the NOTICE | ||
| * file distributed with this work for additional information regarding copyright ownership. The ASF licenses this file | ||
| * to You under the Apache License, Version 2.0 (the "License"); you may not use this file except in compliance with the | ||
| * License. You may obtain a copy of the License at | ||
| * | ||
| * http://www.apache.org/licenses/LICENSE-2.0 | ||
| * | ||
| * Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on | ||
| * an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the | ||
| * specific language governing permissions and limitations under the License. | ||
| */ | ||
|
|
||
| package org.apache.kafka.common.errors; | ||
|
|
||
| public class NotControllerException extends RetriableException { | ||
|
|
||
| private static final long serialVersionUID = 1L; | ||
|
|
||
| public NotControllerException(String message) { | ||
| super(message); | ||
| } | ||
|
|
||
| public NotControllerException(String message, Throwable cause) { | ||
| super(message, cause); | ||
| } | ||
|
|
||
| } |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -770,6 +770,50 @@ public class Protocol { | |
| public static final Schema[] API_VERSIONS_REQUEST = new Schema[]{API_VERSIONS_REQUEST_V0}; | ||
| public static final Schema[] API_VERSIONS_RESPONSE = new Schema[]{API_VERSIONS_RESPONSE_V0}; | ||
|
|
||
| /* Admin requests common */ | ||
| public static final Schema CONFIG_ENTRY = new Schema(new Field("config_key", STRING, "Configuration key name"), | ||
| new Field("config_value", STRING, "Configuration value")); | ||
|
|
||
| public static final Schema PARTITION_REPLICA_ASSIGNMENT_ENTRY = new Schema( | ||
| new Field("partition_id", INT32), | ||
| new Field("replicas", new ArrayOf(INT32), "The set of all nodes that should host this partition. The first replica in the list is the preferred leader.")); | ||
|
|
||
| public static final Schema TOPIC_ERROR_CODE = new Schema(new Field("topic", STRING), new Field("error_code", INT16)); | ||
|
|
||
| /* CreateTopic api */ | ||
| public static final Schema SINGLE_CREATE_TOPIC_REQUEST_V0 = new Schema( | ||
| new Field("topic", | ||
| STRING, | ||
| "Name for newly created topic."), | ||
| new Field("num_partitions", | ||
| INT32, | ||
| "Number of partitions to be created. -1 indicates unset."), | ||
| new Field("replication_factor", | ||
| INT16, | ||
| "Replication factor for the topic. -1 indicates unset."), | ||
| new Field("replica_assignment", | ||
| new ArrayOf(PARTITION_REPLICA_ASSIGNMENT_ENTRY), | ||
| "Replica assignment among kafka brokers for this topic partitions. If this is set num_partitions and replication_factor must be unset."), | ||
| new Field("configs", | ||
| new ArrayOf(CONFIG_ENTRY), | ||
| "Topic level configuration for topic to be set.")); | ||
|
|
||
| public static final Schema CREATE_TOPICS_REQUEST_V0 = new Schema( | ||
| new Field("create_topic_requests", | ||
| new ArrayOf(SINGLE_CREATE_TOPIC_REQUEST_V0), | ||
| "An array of single topic creation requests. Can not have multiple entries for the same topic."), | ||
| new Field("timeout", | ||
| INT32, | ||
| "The time in ms to wait for a topic to be completely created on the controller node. Values <= 0 will trigger topic creation and return immediatly")); | ||
|
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. typo immediatly
Member
Author
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. ack |
||
|
|
||
| public static final Schema CREATE_TOPICS_RESPONSE_V0 = new Schema( | ||
| new Field("topic_error_codes", | ||
| new ArrayOf(TOPIC_ERROR_CODE), | ||
| "An array of per topic error codes.")); | ||
|
|
||
| public static final Schema[] CREATE_TOPICS_REQUEST = new Schema[] {CREATE_TOPICS_REQUEST_V0}; | ||
| public static final Schema[] CREATE_TOPICS_RESPONSE = new Schema[] {CREATE_TOPICS_RESPONSE_V0}; | ||
|
|
||
| /* an array of all requests and responses with all schema versions; a null value in the inner array means that the | ||
| * particular version is not supported */ | ||
| public static final Schema[][] REQUESTS = new Schema[ApiKeys.MAX_API_KEY + 1][]; | ||
|
|
@@ -799,6 +843,7 @@ public class Protocol { | |
| REQUESTS[ApiKeys.LIST_GROUPS.id] = LIST_GROUPS_REQUEST; | ||
| REQUESTS[ApiKeys.SASL_HANDSHAKE.id] = SASL_HANDSHAKE_REQUEST; | ||
| REQUESTS[ApiKeys.API_VERSIONS.id] = API_VERSIONS_REQUEST; | ||
| REQUESTS[ApiKeys.CREATE_TOPICS.id] = CREATE_TOPICS_REQUEST; | ||
|
|
||
| RESPONSES[ApiKeys.PRODUCE.id] = PRODUCE_RESPONSE; | ||
| RESPONSES[ApiKeys.FETCH.id] = FETCH_RESPONSE; | ||
|
|
@@ -819,6 +864,7 @@ public class Protocol { | |
| RESPONSES[ApiKeys.LIST_GROUPS.id] = LIST_GROUPS_RESPONSE; | ||
| RESPONSES[ApiKeys.SASL_HANDSHAKE.id] = SASL_HANDSHAKE_RESPONSE; | ||
| RESPONSES[ApiKeys.API_VERSIONS.id] = API_VERSIONS_RESPONSE; | ||
| RESPONSES[ApiKeys.CREATE_TOPICS.id] = CREATE_TOPICS_RESPONSE; | ||
|
|
||
| /* set the minimum and maximum version of each api */ | ||
| for (ApiKeys api : ApiKeys.values()) { | ||
|
|
||
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
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.
Since we haven't deprecated ErrorMapping, could we add those new error codes as place holders in ErrorMapping?
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.
I will add it.
Note: That since only a newer client (the admin client) can actually make a request that causes this error, it should never show up in the old scala clients.
Uh oh!
There was an error while loading. Please reload this page.
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.
Looking closer at ErrorMapping there are many codes left out that would never be used. Do these fall under that category? Are you sure I should add them?
For now I will add the comment like the others:
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.
Yes, we just want to prevent someone from accidentally adding the same error code for a different exception in ErrorMapping in the future.