-
Notifications
You must be signed in to change notification settings - Fork 13.3k
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
[FLINK-9955] Add Kubernetes ClusterDescriptor to support deploying session cluster. #9973
Conversation
Thanks a lot for your contribution to the Apache Flink project. I'm the @flinkbot. I help the community Automated ChecksLast check on commit bce759c (Wed Dec 04 15:57:18 UTC 2019) Warnings:
Mention the bot in a comment to re-run the automated checks. Review Progress
Please see the Pull Request Review Guide for a full explanation of the review process. The Bot is tracking the review progress through labels. Labels are applied according to the order of the review items. For consensus, approval by a Flink committer of PMC member is required Bot commandsThe @flinkbot bot supports the following commands:
|
58b9eb0
to
450cd10
Compare
CI report:
Bot commandsThe @flinkbot bot supports the following commands:
|
450cd10
to
6285ff3
Compare
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.
Other than the definition of logException. It generally LGTM. +1 for merging.
6148344
to
baba3f1
Compare
1f1ec2e
to
62d628c
Compare
|
||
private static final Logger LOG = LoggerFactory.getLogger(KubernetesClusterDescriptor.class); | ||
|
||
private static final String CLUSTER_DESCRIPTION = "Kubernetes cluster"; |
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.
It could be a follow up we give detailed description as on YARN. Just thought, not requirement for this patch :-) It would be better if we track it as a ticket on JIRA.
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.
Thanks for your suggestion. I have created a ticket FLINK-14986 to track this.
bce759c
to
63bc2b0
Compare
@tisonkun I have rebase the master. Could you please take a look again? |
flinkConfig.setString(KubernetesConfigOptionsInternal.ENTRY_POINT_CLASS, entryPoint); | ||
|
||
// Rpc(6123), blob(6124), rest(8081) taskManagerRpc(6122) port need to be exposed, so update them to fixed port. | ||
flinkConfig.setString(BlobServerOptions.PORT, String.valueOf(Constants.BLOB_SERVER_PORT)); |
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.
So user setups on these PORT
s have no power? It is a strong constraint?
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.
Nice catch.
Fixed port is required, because we need to expose the port in kubernetes. I will update the PR so that it could respect the user specified port value. If no BlobServerOptions.PORT
and TaskManagerOptions.RPC_PORT
are set, the default value will be used.
5bfa111
to
e8267bc
Compare
flinkConfig.setString(KubernetesConfigOptionsInternal.ENTRY_POINT_CLASS, entryPoint); | ||
|
||
// Rpc(6123), blob(6124), rest(8081) taskManagerRpc(6122) port need to be exposed, so update them to fixed port. | ||
if (Integer.valueOf(flinkConfig.get(BlobServerOptions.PORT)) == 0) { |
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.
Possibly we catch the parse exception and throw/log a specific information that port range doesn't support in K8S deployment. And when we config the default port on user configure 0
, it'd better we log a warning so that user knows his original purpose "random choose a port" isn't respected.
{
final int blobServerPort;
try {
blobServerPort = Integer.parseInt(flinkConfig.get(BlobServerOptions.PORT));
} catch (NumberFormatException e) {
// log...
throw new ClusterDeploymentException("...");
}
if (blobServerPort == 0) {
flinkConfig.setString(BlobServerOptions.PORT, String.valueOf(Constants.BLOB_SERVER_PORT));
// log ...
}
}
{
final int taskManagerRpcPort;
try {
taskManagerRpcPort = Integer.parseInt(flinkConfig.get(BlobServerOptions.PORT));
} catch (NumberFormatException e) {
// log...
throw new ClusterDeploymentException("...");
}
if (taskManagerRpcPort == 0) {
flinkConfig.setString(TaskManagerOptions.RPC_PORT, String.valueOf(Constants.TASK_MANAGER_RPC_PORT));
// log ...
}
}
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 a method parsePort
in KubernetesUtils
so that it could be reused by FlinkMasterDeploymentDecorator
and TaskManagerPodDecorator
.
e8267bc
to
b0f46f5
Compare
…ager pod respect user defined config BlobServerOptions.PORT and TaskManagerOptions.RPC_PORT need to be respected.
…deploying session cluster
b0f46f5
to
1bd7fa1
Compare
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.
LGTM. Merging...
What is the purpose of the change
This PR is part of FLINK-9953.
KubernetesClusterDescriptor
is added to support deploy session cluster.This PR is based on #9957 #9965 #9986.
Brief change log
Verifying this change
This change added related unit tests.
Does this pull request potentially affect one of the following parts:
@Public(Evolving)
: (no)Documentation