Aitozi commented on a change in pull request #131:
URL:
https://github.com/apache/flink-kubernetes-operator/pull/131#discussion_r838131907
##########
File path:
flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/FlinkConfigBuilder.java
##########
@@ -50,13 +51,20 @@
/** Builder to get effective flink config from {@link FlinkDeployment}. */
public class FlinkConfigBuilder {
- private final FlinkDeployment deploy;
+ private final ObjectMeta meta;
private final FlinkDeploymentSpec spec;
private final Configuration effectiveConfig;
public FlinkConfigBuilder(FlinkDeployment deploy, Configuration
flinkConfig) {
- this.deploy = deploy;
- this.spec = this.deploy.getSpec();
+ this.meta = deploy.getMetadata();
Review comment:
this(deploy.getMetadata(), deploy.getSpec(), flinkConfig)
##########
File path:
flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/FlinkConfigBuilder.java
##########
@@ -50,13 +51,20 @@
/** Builder to get effective flink config from {@link FlinkDeployment}. */
public class FlinkConfigBuilder {
- private final FlinkDeployment deploy;
+ private final ObjectMeta meta;
private final FlinkDeploymentSpec spec;
private final Configuration effectiveConfig;
public FlinkConfigBuilder(FlinkDeployment deploy, Configuration
flinkConfig) {
- this.deploy = deploy;
- this.spec = this.deploy.getSpec();
+ this.meta = deploy.getMetadata();
+ this.spec = deploy.getSpec();
+ this.effectiveConfig = new Configuration(flinkConfig);
+ }
+
+ public FlinkConfigBuilder(
Review comment:
Can we directly pass the namespace and clusterId here? It seems we only
use these two fields
##########
File path:
flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/FlinkDeploymentController.java
##########
@@ -104,9 +105,30 @@ public DeleteControl cleanup(FlinkDeployment flinkApp,
Context context) {
@Override
public UpdateControl<FlinkDeployment> reconcile(FlinkDeployment flinkApp,
Context context) {
- FlinkDeployment originalCopy = ReconciliationUtils.clone(flinkApp);
LOG.info("Starting reconciliation");
+ FlinkDeployment originalCopy = ReconciliationUtils.clone(flinkApp);
+ FlinkDeploymentSpec lastReconciledSpec =
+
flinkApp.getStatus().getReconciliationStatus().getLastReconciledSpec();
+ if (lastReconciledSpec != null) {
+ Configuration latestValidatedConfig =
+ FlinkUtils.getEffectiveConfig(
+ flinkApp.getMetadata(),
+ lastReconciledSpec,
+ defaultConfig.getFlinkConfig());
+ try {
+ observerFactory
+ .getOrCreate(flinkApp)
+ .observe(flinkApp, context, latestValidatedConfig);
+ } catch (DeploymentFailedException dfe) {
+ handleDeploymentFailed(flinkApp, dfe);
+ LOG.info("Reconciliation successfully completed");
Review comment:
This log is not right
##########
File path:
flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/FlinkDeploymentController.java
##########
@@ -104,9 +105,30 @@ public DeleteControl cleanup(FlinkDeployment flinkApp,
Context context) {
@Override
public UpdateControl<FlinkDeployment> reconcile(FlinkDeployment flinkApp,
Context context) {
- FlinkDeployment originalCopy = ReconciliationUtils.clone(flinkApp);
LOG.info("Starting reconciliation");
+ FlinkDeployment originalCopy = ReconciliationUtils.clone(flinkApp);
+ FlinkDeploymentSpec lastReconciledSpec =
+
flinkApp.getStatus().getReconciliationStatus().getLastReconciledSpec();
+ if (lastReconciledSpec != null) {
+ Configuration latestValidatedConfig =
Review comment:
nit: maybe called `lastValidatedConfig` better, because it from the
`lastReconciledSpec`
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]