From 78cbc19bddfb6da38a00a9fce5ec7c011ecb3c72 Mon Sep 17 00:00:00 2001 From: "Tzu-Li (Gordon) Tai" Date: Fri, 4 Sep 2020 11:43:59 +0800 Subject: [PATCH 1/9] [FLINK-19130] [core] Extend metric interfaces for backpressure metrics This commit extends the FunctionTypeMetrics interface and introduces a new FunctionDispatcherMetrics interface with a goal to expose the following backpressure-related metrics: - Number of blocked addresses (per function type) - Number of inflight async ops (per function type + per-operator) --- .../FlinkFunctionDispatcherMetrics.java | 43 +++++++++++++++++++ .../metrics/FlinkFunctionTypeMetrics.java | 24 +++++++++++ .../metrics/FunctionDispatcherMetrics.java | 26 +++++++++++ .../core/metrics/FunctionTypeMetrics.java | 8 ++++ 4 files changed, 101 insertions(+) create mode 100644 statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FlinkFunctionDispatcherMetrics.java create mode 100644 statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FunctionDispatcherMetrics.java diff --git a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FlinkFunctionDispatcherMetrics.java b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FlinkFunctionDispatcherMetrics.java new file mode 100644 index 000000000..c060a3dd8 --- /dev/null +++ b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FlinkFunctionDispatcherMetrics.java @@ -0,0 +1,43 @@ +/* + * 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.flink.statefun.flink.core.metrics; + +import java.util.Objects; +import org.apache.flink.metrics.Counter; +import org.apache.flink.metrics.MetricGroup; + +public class FlinkFunctionDispatcherMetrics implements FunctionDispatcherMetrics { + private final Counter inflightAsyncOperations; + + public FlinkFunctionDispatcherMetrics(MetricGroup operatorGroup) { + Objects.requireNonNull(operatorGroup, "operatorGroup"); + + this.inflightAsyncOperations = operatorGroup.counter("inflight-async-ops"); + } + + @Override + public void asyncOperationRegistered() { + inflightAsyncOperations.inc(); + } + + @Override + public void asyncOperationCompleted() { + inflightAsyncOperations.dec(); + } +} diff --git a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FlinkFunctionTypeMetrics.java b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FlinkFunctionTypeMetrics.java index d9eec0cac..553cd6845 100644 --- a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FlinkFunctionTypeMetrics.java +++ b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FlinkFunctionTypeMetrics.java @@ -27,12 +27,16 @@ final class FlinkFunctionTypeMetrics implements FunctionTypeMetrics { private final Counter outgoingLocalMessage; private final Counter outgoingRemoteMessage; private final Counter outgoingEgress; + private final Counter blockedAddress; + private final Counter inflightAsyncOps; FlinkFunctionTypeMetrics(MetricGroup typeGroup) { this.incoming = metered(typeGroup, "in"); this.outgoingLocalMessage = metered(typeGroup, "out-local"); this.outgoingRemoteMessage = metered(typeGroup, "out-remote"); this.outgoingEgress = metered(typeGroup, "out-egress"); + this.blockedAddress = typeGroup.counter("num-blocked-address"); + this.inflightAsyncOps = typeGroup.counter("inflight-async-ops"); } @Override @@ -55,6 +59,26 @@ public void outgoingEgressMessage() { this.outgoingEgress.inc(); } + @Override + public void blockedAddress() { + this.blockedAddress.inc(); + } + + @Override + public void unblockedAddress() { + this.blockedAddress.dec(); + } + + @Override + public void asyncOperationRegistered() { + this.inflightAsyncOps.inc(); + } + + @Override + public void asyncOperationCompleted() { + this.inflightAsyncOps.dec(); + } + private static SimpleCounter metered(MetricGroup metrics, String name) { SimpleCounter counter = metrics.counter(name, new SimpleCounter()); metrics.meter(name + "Rate", new MeterView(counter, 60)); diff --git a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FunctionDispatcherMetrics.java b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FunctionDispatcherMetrics.java new file mode 100644 index 000000000..f8d6dc22c --- /dev/null +++ b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FunctionDispatcherMetrics.java @@ -0,0 +1,26 @@ +/* + * 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.flink.statefun.flink.core.metrics; + +public interface FunctionDispatcherMetrics { + + void asyncOperationRegistered(); + + void asyncOperationCompleted(); +} diff --git a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FunctionTypeMetrics.java b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FunctionTypeMetrics.java index 54a84e347..acfc146b7 100644 --- a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FunctionTypeMetrics.java +++ b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FunctionTypeMetrics.java @@ -26,4 +26,12 @@ public interface FunctionTypeMetrics { void outgoingRemoteMessage(); void outgoingEgressMessage(); + + void blockedAddress(); + + void unblockedAddress(); + + void asyncOperationRegistered(); + + void asyncOperationCompleted(); } From 640886603f1778d999dc46ebb2be159df7eb30b5 Mon Sep 17 00:00:00 2001 From: "Tzu-Li (Gordon) Tai" Date: Fri, 4 Sep 2020 11:47:53 +0800 Subject: [PATCH 2/9] [FLINK-19130] [core] Let StatefulFunctionsRepository extend FunctionTypeMetricsRepository This commit introduces a new scoped-down interface FunctionTypeMetricsRepository that the existing StatefulFunctionsRepository now implements. This interface will be provided to components that needs access to per-function metrics. --- .../functions/StatefulFunctionRepository.java | 9 ++++++- .../FunctionTypeMetricsRepository.java | 25 +++++++++++++++++++ 2 files changed, 33 insertions(+), 1 deletion(-) create mode 100644 statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FunctionTypeMetricsRepository.java diff --git a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/functions/StatefulFunctionRepository.java b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/functions/StatefulFunctionRepository.java index 7c1af5174..d2af98273 100644 --- a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/functions/StatefulFunctionRepository.java +++ b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/functions/StatefulFunctionRepository.java @@ -24,13 +24,15 @@ import org.apache.flink.statefun.flink.core.di.Label; import org.apache.flink.statefun.flink.core.message.MessageFactory; import org.apache.flink.statefun.flink.core.metrics.FunctionTypeMetrics; +import org.apache.flink.statefun.flink.core.metrics.FunctionTypeMetricsRepository; import org.apache.flink.statefun.flink.core.metrics.MetricsFactory; import org.apache.flink.statefun.flink.core.state.FlinkStateBinder; import org.apache.flink.statefun.flink.core.state.PersistedStates; import org.apache.flink.statefun.flink.core.state.State; import org.apache.flink.statefun.sdk.FunctionType; -final class StatefulFunctionRepository implements FunctionRepository { +final class StatefulFunctionRepository + implements FunctionRepository, FunctionTypeMetricsRepository { private final ObjectOpenHashMap instances; private final State flinkState; private final FunctionLoader functionLoader; @@ -59,6 +61,11 @@ public LiveFunction get(FunctionType type) { return function; } + @Override + public FunctionTypeMetrics getMetrics(FunctionType functionType) { + return get(functionType).metrics(); + } + private StatefulFunction load(FunctionType functionType) { org.apache.flink.statefun.sdk.StatefulFunction statefulFunction = functionLoader.load(functionType); diff --git a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FunctionTypeMetricsRepository.java b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FunctionTypeMetricsRepository.java new file mode 100644 index 000000000..94a9a2f70 --- /dev/null +++ b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FunctionTypeMetricsRepository.java @@ -0,0 +1,25 @@ +/* + * 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.flink.statefun.flink.core.metrics; + +import org.apache.flink.statefun.sdk.FunctionType; + +public interface FunctionTypeMetricsRepository { + FunctionTypeMetrics getMetrics(FunctionType functionType); +} From 0fbb17afbecde697f8a099aa758dba45a1a95f7b Mon Sep 17 00:00:00 2001 From: "Tzu-Li (Gordon) Tai" Date: Fri, 4 Sep 2020 11:49:43 +0800 Subject: [PATCH 3/9] [FLINK-19130] [core] Allow ObjectContainer to add alias keys This extends the ObjectContainer to set alias keys that returns the same instance. This is required for the StatefulFunctionsRepository, since we want to share the same instance for different keys (depending on which interface we expose to different components). --- .../flink/core/di/ObjectContainer.java | 8 ++++ .../flink/core/di/ObjectContainerTest.java | 44 +++++++++++++++++++ 2 files changed, 52 insertions(+) create mode 100644 statefun-flink/statefun-flink-core/src/test/java/org/apache/flink/statefun/flink/core/di/ObjectContainerTest.java diff --git a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/di/ObjectContainer.java b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/di/ObjectContainer.java index 7292a2e22..8d43bcf2c 100644 --- a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/di/ObjectContainer.java +++ b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/di/ObjectContainer.java @@ -51,6 +51,14 @@ public void add(String label, Class type, Class actual) { factories.put(new Key(type, label), () -> createReflectively(actual)); } + public void addAlias( + String newLabel, + Class newType, + String existingLabel, + Class existingType) { + factories.put(new Key(newType, newLabel), () -> get(existingType, existingLabel)); + } + public void add(String label, Lazy lazyValue) { factories.put(new Key(Lazy.class, label), () -> lazyValue.withContainer(this)); } diff --git a/statefun-flink/statefun-flink-core/src/test/java/org/apache/flink/statefun/flink/core/di/ObjectContainerTest.java b/statefun-flink/statefun-flink-core/src/test/java/org/apache/flink/statefun/flink/core/di/ObjectContainerTest.java new file mode 100644 index 000000000..b532fedfa --- /dev/null +++ b/statefun-flink/statefun-flink-core/src/test/java/org/apache/flink/statefun/flink/core/di/ObjectContainerTest.java @@ -0,0 +1,44 @@ +/* + * 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.flink.statefun.flink.core.di; + +import static org.hamcrest.CoreMatchers.theInstance; +import static org.hamcrest.MatcherAssert.assertThat; + +import org.junit.Test; + +public class ObjectContainerTest { + + @Test + public void addAliasTest() { + final ObjectContainer container = new ObjectContainer(); + + container.add("label-1", InterfaceA.class, TestClass.class); + container.addAlias("label-2", InterfaceB.class, "label-1", InterfaceA.class); + + assertThat( + container.get(InterfaceB.class, "label-2"), + theInstance(container.get(InterfaceA.class, "label-1"))); + } + + private interface InterfaceA {} + + private interface InterfaceB {} + + private static class TestClass implements InterfaceA, InterfaceB {} +} From 6c77ca48883ad87eebc0899bc2ecddee18d14092 Mon Sep 17 00:00:00 2001 From: "Tzu-Li (Gordon) Tai" Date: Fri, 4 Sep 2020 11:51:44 +0800 Subject: [PATCH 4/9] [FLINK-19130] [core] Wire-in new metrics into AsyncSink --- .../core/backpressure/BackPressureValve.java | 8 +++++ .../ThresholdBackPressureValve.java | 6 ++++ .../flink/core/functions/AsyncSink.java | 31 ++++++++++++++++--- .../flink/core/functions/Reductions.java | 12 +++++++ .../flink/core/functions/ReusableContext.java | 2 +- .../flink/core/functions/ReductionsTest.java | 3 +- 6 files changed, 56 insertions(+), 6 deletions(-) diff --git a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/backpressure/BackPressureValve.java b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/backpressure/BackPressureValve.java index 29f1c1945..393f41b80 100644 --- a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/backpressure/BackPressureValve.java +++ b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/backpressure/BackPressureValve.java @@ -53,4 +53,12 @@ public interface BackPressureValve { * @param address the address */ void blockAddress(Address address); + + /** + * Checks whether a given address was previously blocked with {@link #blockAddress(Address)}. + * + * @param address the address to check + * @return boolean indicating whether or not the address was blocked. + */ + boolean isAddressBlocked(Address address); } diff --git a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/backpressure/ThresholdBackPressureValve.java b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/backpressure/ThresholdBackPressureValve.java index 1ce0b007d..91b642b24 100644 --- a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/backpressure/ThresholdBackPressureValve.java +++ b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/backpressure/ThresholdBackPressureValve.java @@ -88,6 +88,12 @@ public void notifyAsyncOperationCompleted(Address owningAddress) { blockedAddressSet.remove(owningAddress); } + /** {@inheritDoc} */ + @Override + public boolean isAddressBlocked(Address address) { + return blockedAddressSet.containsKey(address); + } + private boolean totalPendingAsyncOperationsAtCapacity() { return maximumPendingAsynchronousOperations > 0 && pendingAsynchronousOperationsCount >= maximumPendingAsynchronousOperations; diff --git a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/functions/AsyncSink.java b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/functions/AsyncSink.java index 873ab3452..647c9683e 100644 --- a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/functions/AsyncSink.java +++ b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/functions/AsyncSink.java @@ -27,6 +27,9 @@ import org.apache.flink.statefun.flink.core.di.Label; import org.apache.flink.statefun.flink.core.di.Lazy; import org.apache.flink.statefun.flink.core.message.Message; +import org.apache.flink.statefun.flink.core.metrics.FunctionDispatcherMetrics; +import org.apache.flink.statefun.flink.core.metrics.FunctionTypeMetrics; +import org.apache.flink.statefun.flink.core.metrics.FunctionTypeMetricsRepository; import org.apache.flink.statefun.flink.core.queue.Locks; import org.apache.flink.statefun.flink.core.queue.MpscQueue; import org.apache.flink.statefun.sdk.Address; @@ -36,6 +39,8 @@ final class AsyncSink { private final Lazy reductions; private final Executor operatorMailbox; private final BackPressureValve backPressureValve; + private final FunctionTypeMetricsRepository metricsRepository; + private final FunctionDispatcherMetrics dispatcherMetrics; private final MpscQueue completed = new MpscQueue<>(32768, Locks.jdkReentrantLock()); @@ -44,14 +49,18 @@ final class AsyncSink { PendingAsyncOperations pendingAsyncOperations, @Label("mailbox-executor") Executor operatorMailbox, @Label("reductions") Lazy reductions, - @Label("backpressure-valve") BackPressureValve backPressureValve) { + @Label("backpressure-valve") BackPressureValve backPressureValve, + @Label("function-metrics-repository") FunctionTypeMetricsRepository metricsRepository, + @Label("function-dispatcher-metrics") FunctionDispatcherMetrics dispatcherMetrics) { this.pendingAsyncOperations = Objects.requireNonNull(pendingAsyncOperations); this.reductions = Objects.requireNonNull(reductions); this.operatorMailbox = Objects.requireNonNull(operatorMailbox); this.backPressureValve = Objects.requireNonNull(backPressureValve); + this.metricsRepository = Objects.requireNonNull(metricsRepository); + this.dispatcherMetrics = Objects.requireNonNull(dispatcherMetrics); } - void accept(Message metadata, CompletableFuture future) { + void accept(Address sourceAddress, Message metadata, CompletableFuture future) { final long futureId = ThreadLocalRandom.current().nextLong(); // TODO: is this is good enough? // we keep the message in state (associated with futureId) until either: // 1. the future successfully completes and the message is processed. The state would be @@ -59,8 +68,11 @@ void accept(Message metadata, CompletableFuture future) { // 2. after recovery, we clear that state by notifying the owning function that we don't know // what happened // with that particular async operation. - pendingAsyncOperations.add(metadata.source(), futureId, metadata); + pendingAsyncOperations.add(sourceAddress, futureId, metadata); backPressureValve.notifyAsyncOperationRegistered(); + + metricsRepository.getMetrics(sourceAddress.type()).asyncOperationRegistered(); + dispatcherMetrics.asyncOperationRegistered(); future.whenComplete((result, throwable) -> enqueue(metadata, futureId, result, throwable)); } @@ -72,6 +84,7 @@ void accept(Message metadata, CompletableFuture future) { */ void blockAddress(Address address) { backPressureValve.blockAddress(address); + metricsRepository.getMetrics(address.type()).blockedAddress(); } private void enqueue(Message message, long futureId, T result, Throwable throwable) { @@ -90,7 +103,17 @@ private void drainOnOperatorThread() { Reductions reductions = this.reductions.get(); Message message; while ((message = batchOfCompletedFutures.poll()) != null) { - backPressureValve.notifyAsyncOperationCompleted(message.target()); + Address target = message.target(); + FunctionTypeMetrics functionMetrics = metricsRepository.getMetrics(target.type()); + + // must check whether address was blocked BEFORE notifying completion + if (backPressureValve.isAddressBlocked(target)) { + functionMetrics.unblockedAddress(); + } + backPressureValve.notifyAsyncOperationCompleted(target); + + functionMetrics.asyncOperationCompleted(); + dispatcherMetrics.asyncOperationCompleted(); reductions.enqueue(message); } reductions.processEnvelopes(); diff --git a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/functions/Reductions.java b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/functions/Reductions.java index abce51afe..44686e012 100644 --- a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/functions/Reductions.java +++ b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/functions/Reductions.java @@ -32,7 +32,10 @@ import org.apache.flink.statefun.flink.core.di.ObjectContainer; import org.apache.flink.statefun.flink.core.message.Message; import org.apache.flink.statefun.flink.core.message.MessageFactory; +import org.apache.flink.statefun.flink.core.metrics.FlinkFunctionDispatcherMetrics; import org.apache.flink.statefun.flink.core.metrics.FlinkMetricsFactory; +import org.apache.flink.statefun.flink.core.metrics.FunctionDispatcherMetrics; +import org.apache.flink.statefun.flink.core.metrics.FunctionTypeMetricsRepository; import org.apache.flink.statefun.flink.core.metrics.MetricsFactory; import org.apache.flink.statefun.flink.core.state.FlinkState; import org.apache.flink.statefun.flink.core.state.State; @@ -71,6 +74,11 @@ static Reductions create( container.add("function-providers", Map.class, statefulFunctionsUniverse.functions()); container.add( "function-repository", FunctionRepository.class, StatefulFunctionRepository.class); + container.addAlias( + "function-metrics-repository", + FunctionTypeMetricsRepository.class, + "function-repository", + FunctionRepository.class); // for FlinkState container.add("runtime-context", RuntimeContext.class, context); @@ -96,6 +104,10 @@ static Reductions create( container.add(Reductions.class); container.add(LocalFunctionGroup.class); container.add("metrics-factory", MetricsFactory.class, new FlinkMetricsFactory(metricGroup)); + container.add( + "function-dispatcher-metrics", + FunctionDispatcherMetrics.class, + new FlinkFunctionDispatcherMetrics(metricGroup)); // for delayed messages container.add( diff --git a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/functions/ReusableContext.java b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/functions/ReusableContext.java index 55c646be6..62289e521 100644 --- a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/functions/ReusableContext.java +++ b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/functions/ReusableContext.java @@ -113,7 +113,7 @@ public void registerAsyncOperation(M metadata, CompletableFuture futur Objects.requireNonNull(future); Message message = messageFactory.from(self(), self(), metadata); - asyncSink.accept(message, future); + asyncSink.accept(self(), message, future); } @Override diff --git a/statefun-flink/statefun-flink-core/src/test/java/org/apache/flink/statefun/flink/core/functions/ReductionsTest.java b/statefun-flink/statefun-flink-core/src/test/java/org/apache/flink/statefun/flink/core/functions/ReductionsTest.java index 46d9ed1e8..bdb67e3a8 100644 --- a/statefun-flink/statefun-flink-core/src/test/java/org/apache/flink/statefun/flink/core/functions/ReductionsTest.java +++ b/statefun-flink/statefun-flink-core/src/test/java/org/apache/flink/statefun/flink/core/functions/ReductionsTest.java @@ -58,6 +58,7 @@ import org.apache.flink.metrics.Gauge; import org.apache.flink.metrics.Meter; import org.apache.flink.metrics.MetricGroup; +import org.apache.flink.metrics.SimpleCounter; import org.apache.flink.runtime.state.KeyGroupedInternalPriorityQueue; import org.apache.flink.runtime.state.Keyed; import org.apache.flink.runtime.state.KeyedStateBackend; @@ -578,7 +579,7 @@ public Counter counter(int i) { @Override public Counter counter(String s) { - throw new UnsupportedOperationException(); + return new SimpleCounter(); } @Override From 6bb5cde09ff4816313d1d81acbe8e74528319fee Mon Sep 17 00:00:00 2001 From: "Tzu-Li (Gordon) Tai" Date: Fri, 4 Sep 2020 11:54:46 +0800 Subject: [PATCH 5/9] [FLINK-19130] [core] Rename MetricsFactory to FunctionTypeMetricsFactory Now that we have multiple types of metrics with different scopes (per-function / per-operator), this interface should be renamed to convey that it is a factory specific for per-function scoped metrics. --- .../flink/statefun/flink/core/functions/Reductions.java | 9 ++++++--- .../flink/core/functions/StatefulFunctionRepository.java | 8 ++++---- ...sFactory.java => FlinkFuncionTypeMetricsFactory.java} | 4 ++-- ...etricsFactory.java => FuncionTypeMetricsFactory.java} | 2 +- 4 files changed, 13 insertions(+), 10 deletions(-) rename statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/{FlinkMetricsFactory.java => FlinkFuncionTypeMetricsFactory.java} (90%) rename statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/{MetricsFactory.java => FuncionTypeMetricsFactory.java} (95%) diff --git a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/functions/Reductions.java b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/functions/Reductions.java index 44686e012..4a90a204e 100644 --- a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/functions/Reductions.java +++ b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/functions/Reductions.java @@ -32,11 +32,11 @@ import org.apache.flink.statefun.flink.core.di.ObjectContainer; import org.apache.flink.statefun.flink.core.message.Message; import org.apache.flink.statefun.flink.core.message.MessageFactory; +import org.apache.flink.statefun.flink.core.metrics.FlinkFuncionTypeMetricsFactory; import org.apache.flink.statefun.flink.core.metrics.FlinkFunctionDispatcherMetrics; -import org.apache.flink.statefun.flink.core.metrics.FlinkMetricsFactory; +import org.apache.flink.statefun.flink.core.metrics.FuncionTypeMetricsFactory; import org.apache.flink.statefun.flink.core.metrics.FunctionDispatcherMetrics; import org.apache.flink.statefun.flink.core.metrics.FunctionTypeMetricsRepository; -import org.apache.flink.statefun.flink.core.metrics.MetricsFactory; import org.apache.flink.statefun.flink.core.state.FlinkState; import org.apache.flink.statefun.flink.core.state.State; import org.apache.flink.statefun.flink.core.types.DynamicallyRegisteredTypes; @@ -103,7 +103,10 @@ static Reductions create( container.add("function-loader", FunctionLoader.class, PredefinedFunctionLoader.class); container.add(Reductions.class); container.add(LocalFunctionGroup.class); - container.add("metrics-factory", MetricsFactory.class, new FlinkMetricsFactory(metricGroup)); + container.add( + "function-metrics-factory", + FuncionTypeMetricsFactory.class, + new FlinkFuncionTypeMetricsFactory(metricGroup)); container.add( "function-dispatcher-metrics", FunctionDispatcherMetrics.class, diff --git a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/functions/StatefulFunctionRepository.java b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/functions/StatefulFunctionRepository.java index d2af98273..5167c8280 100644 --- a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/functions/StatefulFunctionRepository.java +++ b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/functions/StatefulFunctionRepository.java @@ -23,9 +23,9 @@ import org.apache.flink.statefun.flink.core.di.Inject; import org.apache.flink.statefun.flink.core.di.Label; import org.apache.flink.statefun.flink.core.message.MessageFactory; +import org.apache.flink.statefun.flink.core.metrics.FuncionTypeMetricsFactory; import org.apache.flink.statefun.flink.core.metrics.FunctionTypeMetrics; import org.apache.flink.statefun.flink.core.metrics.FunctionTypeMetricsRepository; -import org.apache.flink.statefun.flink.core.metrics.MetricsFactory; import org.apache.flink.statefun.flink.core.state.FlinkStateBinder; import org.apache.flink.statefun.flink.core.state.PersistedStates; import org.apache.flink.statefun.flink.core.state.State; @@ -36,18 +36,18 @@ final class StatefulFunctionRepository private final ObjectOpenHashMap instances; private final State flinkState; private final FunctionLoader functionLoader; - private final MetricsFactory metricsFactory; + private final FuncionTypeMetricsFactory metricsFactory; private final MessageFactory messageFactory; @Inject StatefulFunctionRepository( @Label("function-loader") FunctionLoader functionLoader, - @Label("metrics-factory") MetricsFactory metricsFactory, + @Label("function-metrics-factory") FuncionTypeMetricsFactory functionMetricsFactory, @Label("state") State state, MessageFactory messageFactory) { this.instances = new ObjectOpenHashMap<>(); this.functionLoader = Objects.requireNonNull(functionLoader); - this.metricsFactory = Objects.requireNonNull(metricsFactory); + this.metricsFactory = Objects.requireNonNull(functionMetricsFactory); this.flinkState = Objects.requireNonNull(state); this.messageFactory = Objects.requireNonNull(messageFactory); } diff --git a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FlinkMetricsFactory.java b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FlinkFuncionTypeMetricsFactory.java similarity index 90% rename from statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FlinkMetricsFactory.java rename to statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FlinkFuncionTypeMetricsFactory.java index d65c896b3..4d003bbbd 100644 --- a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FlinkMetricsFactory.java +++ b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FlinkFuncionTypeMetricsFactory.java @@ -21,11 +21,11 @@ import org.apache.flink.metrics.MetricGroup; import org.apache.flink.statefun.sdk.FunctionType; -public class FlinkMetricsFactory implements MetricsFactory { +public class FlinkFuncionTypeMetricsFactory implements FuncionTypeMetricsFactory { private final MetricGroup metricGroup; - public FlinkMetricsFactory(MetricGroup metricGroup) { + public FlinkFuncionTypeMetricsFactory(MetricGroup metricGroup) { this.metricGroup = Objects.requireNonNull(metricGroup); } diff --git a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/MetricsFactory.java b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FuncionTypeMetricsFactory.java similarity index 95% rename from statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/MetricsFactory.java rename to statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FuncionTypeMetricsFactory.java index 7b558f319..2634ed096 100644 --- a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/MetricsFactory.java +++ b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FuncionTypeMetricsFactory.java @@ -19,7 +19,7 @@ import org.apache.flink.statefun.sdk.FunctionType; -public interface MetricsFactory { +public interface FuncionTypeMetricsFactory { FunctionTypeMetrics forType(FunctionType functionType); } From 7ed411d6d50221a94b6917ca55dd90a694e603e7 Mon Sep 17 00:00:00 2001 From: "Tzu-Li (Gordon) Tai" Date: Tue, 8 Sep 2020 13:45:34 +0800 Subject: [PATCH 6/9] [FLINK-19130] [core] Rename AsyncWaiter to InternalContext As a preparation to expose FunctionTypeMetrics to functions for internal axcess only, AsyncWaiter is renamed to a more general-purpose name "InternalContext" so that it makes sense to add more internal-only context methods there. --- .../backpressure/{AsyncWaiter.java => InternalContext.java} | 3 ++- .../flink/core/backpressure/ThresholdBackPressureValve.java | 2 +- .../flink/statefun/flink/core/functions/ReusableContext.java | 4 ++-- .../statefun/flink/core/reqreply/RequestReplyFunction.java | 4 ++-- .../flink/core/reqreply/RequestReplyFunctionTest.java | 5 ++--- 5 files changed, 9 insertions(+), 9 deletions(-) rename statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/backpressure/{AsyncWaiter.java => InternalContext.java} (94%) diff --git a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/backpressure/AsyncWaiter.java b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/backpressure/InternalContext.java similarity index 94% rename from statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/backpressure/AsyncWaiter.java rename to statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/backpressure/InternalContext.java index ddbe247e2..a40d7c83c 100644 --- a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/backpressure/AsyncWaiter.java +++ b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/backpressure/InternalContext.java @@ -19,9 +19,10 @@ package org.apache.flink.statefun.flink.core.backpressure; import org.apache.flink.annotation.Internal; +import org.apache.flink.statefun.sdk.Context; @Internal -public interface AsyncWaiter { +public interface InternalContext extends Context { /** * Signals the runtime to stop invoking the currently executing function with new input until at diff --git a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/backpressure/ThresholdBackPressureValve.java b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/backpressure/ThresholdBackPressureValve.java index 91b642b24..e6f6ade26 100644 --- a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/backpressure/ThresholdBackPressureValve.java +++ b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/backpressure/ThresholdBackPressureValve.java @@ -43,7 +43,7 @@ public final class ThresholdBackPressureValve implements BackPressureValve { /** * a set of address that had explicitly requested to stop processing any new inputs (via {@link - * AsyncWaiter#awaitAsyncOperationComplete()}. Note that this is a set implemented on top of a + * InternalContext#awaitAsyncOperationComplete()}. Note that this is a set implemented on top of a * map, and the value (Boolean) has no meaning. */ private final ObjectOpenHashMap blockedAddressSet = diff --git a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/functions/ReusableContext.java b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/functions/ReusableContext.java index 62289e521..97561dd45 100644 --- a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/functions/ReusableContext.java +++ b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/functions/ReusableContext.java @@ -20,7 +20,7 @@ import java.time.Duration; import java.util.Objects; import java.util.concurrent.CompletableFuture; -import org.apache.flink.statefun.flink.core.backpressure.AsyncWaiter; +import org.apache.flink.statefun.flink.core.backpressure.InternalContext; import org.apache.flink.statefun.flink.core.di.Inject; import org.apache.flink.statefun.flink.core.di.Label; import org.apache.flink.statefun.flink.core.message.Message; @@ -29,7 +29,7 @@ import org.apache.flink.statefun.sdk.Address; import org.apache.flink.statefun.sdk.io.EgressIdentifier; -final class ReusableContext implements ApplyingContext, AsyncWaiter { +final class ReusableContext implements ApplyingContext, InternalContext { private final Partition thisPartition; private final LocalSink localSink; private final RemoteSink remoteSink; diff --git a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/reqreply/RequestReplyFunction.java b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/reqreply/RequestReplyFunction.java index 4defa30ad..bac18c455 100644 --- a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/reqreply/RequestReplyFunction.java +++ b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/reqreply/RequestReplyFunction.java @@ -26,7 +26,7 @@ import java.time.Duration; import java.util.Objects; import java.util.concurrent.CompletableFuture; -import org.apache.flink.statefun.flink.core.backpressure.AsyncWaiter; +import org.apache.flink.statefun.flink.core.backpressure.InternalContext; import org.apache.flink.statefun.flink.core.polyglot.generated.FromFunction; import org.apache.flink.statefun.flink.core.polyglot.generated.FromFunction.EgressMessage; import org.apache.flink.statefun.flink.core.polyglot.generated.FromFunction.InvocationResponse; @@ -111,7 +111,7 @@ private void onRequest(Context context, Any message) { // we need to signal to the runtime that we are unable to process any new input // and we must wait for our in flight asynchronous operation to complete before // we are able to process more input. - ((AsyncWaiter) context).awaitAsyncOperationComplete(); + ((InternalContext) context).awaitAsyncOperationComplete(); } } diff --git a/statefun-flink/statefun-flink-core/src/test/java/org/apache/flink/statefun/flink/core/reqreply/RequestReplyFunctionTest.java b/statefun-flink/statefun-flink-core/src/test/java/org/apache/flink/statefun/flink/core/reqreply/RequestReplyFunctionTest.java index bdb10a968..339cd91eb 100644 --- a/statefun-flink/statefun-flink-core/src/test/java/org/apache/flink/statefun/flink/core/reqreply/RequestReplyFunctionTest.java +++ b/statefun-flink/statefun-flink-core/src/test/java/org/apache/flink/statefun/flink/core/reqreply/RequestReplyFunctionTest.java @@ -36,7 +36,7 @@ import java.util.concurrent.CompletableFuture; import java.util.function.Supplier; import org.apache.flink.statefun.flink.core.TestUtils; -import org.apache.flink.statefun.flink.core.backpressure.AsyncWaiter; +import org.apache.flink.statefun.flink.core.backpressure.InternalContext; import org.apache.flink.statefun.flink.core.httpfn.StateSpec; import org.apache.flink.statefun.flink.core.polyglot.generated.FromFunction; import org.apache.flink.statefun.flink.core.polyglot.generated.FromFunction.DelayedInvocation; @@ -49,7 +49,6 @@ import org.apache.flink.statefun.sdk.Address; import org.apache.flink.statefun.sdk.AsyncOperationResult; import org.apache.flink.statefun.sdk.AsyncOperationResult.Status; -import org.apache.flink.statefun.sdk.Context; import org.apache.flink.statefun.sdk.FunctionType; import org.apache.flink.statefun.sdk.io.EgressIdentifier; import org.junit.Test; @@ -249,7 +248,7 @@ public int capturedInvocationBatchSize() { } } - private static final class FakeContext implements Context, AsyncWaiter { + private static final class FakeContext implements InternalContext { Address caller; boolean needsWaiting; From de60b6456b7bbfc42aef0a6ebacf726eb0c8c69a Mon Sep 17 00:00:00 2001 From: "Tzu-Li (Gordon) Tai" Date: Tue, 8 Sep 2020 13:50:34 +0800 Subject: [PATCH 7/9] [FLINK-19130] [core] Expose function type metrics via InternalContext --- .../core/backpressure/InternalContext.java | 8 +++++ .../flink/core/functions/ReusableContext.java | 6 ++++ .../reqreply/RequestReplyFunctionTest.java | 35 +++++++++++++++++++ 3 files changed, 49 insertions(+) diff --git a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/backpressure/InternalContext.java b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/backpressure/InternalContext.java index a40d7c83c..7ac680366 100644 --- a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/backpressure/InternalContext.java +++ b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/backpressure/InternalContext.java @@ -19,6 +19,7 @@ package org.apache.flink.statefun.flink.core.backpressure; import org.apache.flink.annotation.Internal; +import org.apache.flink.statefun.flink.core.metrics.FunctionTypeMetrics; import org.apache.flink.statefun.sdk.Context; @Internal @@ -37,4 +38,11 @@ public interface InternalContext extends Context { * every async operation registered per each address. */ void awaitAsyncOperationComplete(); + + /** + * Returns the metrics handle for the current invoked function's type. + * + * @return the metrics handle for the current invoked function's type. + */ + FunctionTypeMetrics functionTypeMetrics(); } diff --git a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/functions/ReusableContext.java b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/functions/ReusableContext.java index 97561dd45..77db7dca9 100644 --- a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/functions/ReusableContext.java +++ b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/functions/ReusableContext.java @@ -25,6 +25,7 @@ import org.apache.flink.statefun.flink.core.di.Label; import org.apache.flink.statefun.flink.core.message.Message; import org.apache.flink.statefun.flink.core.message.MessageFactory; +import org.apache.flink.statefun.flink.core.metrics.FunctionTypeMetrics; import org.apache.flink.statefun.flink.core.state.State; import org.apache.flink.statefun.sdk.Address; import org.apache.flink.statefun.sdk.io.EgressIdentifier; @@ -121,6 +122,11 @@ public void awaitAsyncOperationComplete() { asyncSink.blockAddress(self()); } + @Override + public FunctionTypeMetrics functionTypeMetrics() { + return function.metrics(); + } + @Override public Address caller() { return in.source(); diff --git a/statefun-flink/statefun-flink-core/src/test/java/org/apache/flink/statefun/flink/core/reqreply/RequestReplyFunctionTest.java b/statefun-flink/statefun-flink-core/src/test/java/org/apache/flink/statefun/flink/core/reqreply/RequestReplyFunctionTest.java index 339cd91eb..1199ab798 100644 --- a/statefun-flink/statefun-flink-core/src/test/java/org/apache/flink/statefun/flink/core/reqreply/RequestReplyFunctionTest.java +++ b/statefun-flink/statefun-flink-core/src/test/java/org/apache/flink/statefun/flink/core/reqreply/RequestReplyFunctionTest.java @@ -38,6 +38,7 @@ import org.apache.flink.statefun.flink.core.TestUtils; import org.apache.flink.statefun.flink.core.backpressure.InternalContext; import org.apache.flink.statefun.flink.core.httpfn.StateSpec; +import org.apache.flink.statefun.flink.core.metrics.FunctionTypeMetrics; import org.apache.flink.statefun.flink.core.polyglot.generated.FromFunction; import org.apache.flink.statefun.flink.core.polyglot.generated.FromFunction.DelayedInvocation; import org.apache.flink.statefun.flink.core.polyglot.generated.FromFunction.EgressMessage; @@ -250,6 +251,8 @@ public int capturedInvocationBatchSize() { private static final class FakeContext implements InternalContext { + private final FunctionTypeMetrics fakeMetrics = new FakeMetrics(); + Address caller; boolean needsWaiting; @@ -262,6 +265,11 @@ public void awaitAsyncOperationComplete() { needsWaiting = true; } + @Override + public FunctionTypeMetrics functionTypeMetrics() { + return fakeMetrics; + } + @Override public Address self() { return new Address(FN_TYPE, "0"); @@ -288,4 +296,31 @@ public void sendAfter(Duration delay, Address to, Object message) { @Override public void registerAsyncOperation(M metadata, CompletableFuture future) {} } + + private static final class FakeMetrics implements FunctionTypeMetrics { + + @Override + public void asyncOperationRegistered() {} + + @Override + public void asyncOperationCompleted() {} + + @Override + public void incomingMessage() {} + + @Override + public void outgoingRemoteMessage() {} + + @Override + public void outgoingEgressMessage() {} + + @Override + public void outgoingLocalMessage() {} + + @Override + public void blockedAddress() {} + + @Override + public void unblockedAddress() {} + } } From eb63dc5095d89574413884fbcf3e7f4940550d9b Mon Sep 17 00:00:00 2001 From: "Tzu-Li (Gordon) Tai" Date: Tue, 8 Sep 2020 13:56:24 +0800 Subject: [PATCH 8/9] [FLINK-19130] [core] Add backlog messages to FunctionTypeMetrics --- .../flink/core/metrics/FlinkFunctionTypeMetrics.java | 12 ++++++++++++ .../flink/core/metrics/FunctionTypeMetrics.java | 4 ++++ .../core/reqreply/RequestReplyFunctionTest.java | 6 ++++++ 3 files changed, 22 insertions(+) diff --git a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FlinkFunctionTypeMetrics.java b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FlinkFunctionTypeMetrics.java index 553cd6845..34fff35cb 100644 --- a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FlinkFunctionTypeMetrics.java +++ b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FlinkFunctionTypeMetrics.java @@ -29,6 +29,7 @@ final class FlinkFunctionTypeMetrics implements FunctionTypeMetrics { private final Counter outgoingEgress; private final Counter blockedAddress; private final Counter inflightAsyncOps; + private final Counter backlogMessage; FlinkFunctionTypeMetrics(MetricGroup typeGroup) { this.incoming = metered(typeGroup, "in"); @@ -37,6 +38,7 @@ final class FlinkFunctionTypeMetrics implements FunctionTypeMetrics { this.outgoingEgress = metered(typeGroup, "out-egress"); this.blockedAddress = typeGroup.counter("num-blocked-address"); this.inflightAsyncOps = typeGroup.counter("inflight-async-ops"); + this.backlogMessage = typeGroup.counter("num-backlog"); } @Override @@ -79,6 +81,16 @@ public void asyncOperationCompleted() { this.inflightAsyncOps.dec(); } + @Override + public void appendBacklogMessages(int count) { + backlogMessage.inc(count); + } + + @Override + public void consumeBacklogMessages(int count) { + backlogMessage.dec(count); + } + private static SimpleCounter metered(MetricGroup metrics, String name) { SimpleCounter counter = metrics.counter(name, new SimpleCounter()); metrics.meter(name + "Rate", new MeterView(counter, 60)); diff --git a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FunctionTypeMetrics.java b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FunctionTypeMetrics.java index acfc146b7..9fd779ae2 100644 --- a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FunctionTypeMetrics.java +++ b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/metrics/FunctionTypeMetrics.java @@ -34,4 +34,8 @@ public interface FunctionTypeMetrics { void asyncOperationRegistered(); void asyncOperationCompleted(); + + void appendBacklogMessages(int count); + + void consumeBacklogMessages(int count); } diff --git a/statefun-flink/statefun-flink-core/src/test/java/org/apache/flink/statefun/flink/core/reqreply/RequestReplyFunctionTest.java b/statefun-flink/statefun-flink-core/src/test/java/org/apache/flink/statefun/flink/core/reqreply/RequestReplyFunctionTest.java index 1199ab798..cd705888b 100644 --- a/statefun-flink/statefun-flink-core/src/test/java/org/apache/flink/statefun/flink/core/reqreply/RequestReplyFunctionTest.java +++ b/statefun-flink/statefun-flink-core/src/test/java/org/apache/flink/statefun/flink/core/reqreply/RequestReplyFunctionTest.java @@ -322,5 +322,11 @@ public void blockedAddress() {} @Override public void unblockedAddress() {} + + @Override + public void appendBacklogMessages(int count) {} + + @Override + public void consumeBacklogMessages(int count) {} } } From c45f94ec208a26bb0329f62cfaf15737dd92bf2b Mon Sep 17 00:00:00 2001 From: "Tzu-Li (Gordon) Tai" Date: Tue, 8 Sep 2020 14:19:01 +0800 Subject: [PATCH 9/9] [FLINK-19130] [core] Apply backlog metrics in RequestReplyFunction --- .../core/reqreply/RequestReplyFunction.java | 19 ++++--- .../reqreply/RequestReplyFunctionTest.java | 54 +++++++++++++++---- 2 files changed, 56 insertions(+), 17 deletions(-) diff --git a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/reqreply/RequestReplyFunction.java b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/reqreply/RequestReplyFunction.java index bac18c455..a254d4391 100644 --- a/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/reqreply/RequestReplyFunction.java +++ b/statefun-flink/statefun-flink-core/src/main/java/org/apache/flink/statefun/flink/core/reqreply/RequestReplyFunction.java @@ -79,17 +79,18 @@ public RequestReplyFunction( @Override public void invoke(Context context, Object input) { + InternalContext castedContext = (InternalContext) context; if (!(input instanceof AsyncOperationResult)) { - onRequest(context, (Any) input); + onRequest(castedContext, (Any) input); return; } @SuppressWarnings("unchecked") AsyncOperationResult result = (AsyncOperationResult) input; - onAsyncResult(context, result); + onAsyncResult(castedContext, result); } - private void onRequest(Context context, Any message) { + private void onRequest(InternalContext context, Any message) { Invocation.Builder invocationBuilder = singeInvocationBuilder(context, message); int inflightOrBatched = requestState.getOrDefault(-1); if (inflightOrBatched < 0) { @@ -106,17 +107,18 @@ private void onRequest(Context context, Any message) { batch.append(invocationBuilder.build()); inflightOrBatched++; requestState.set(inflightOrBatched); + context.functionTypeMetrics().appendBacklogMessages(1); if (isMaxNumBatchRequestsExceeded(inflightOrBatched)) { // we are at capacity, can't add anything to the batch. // we need to signal to the runtime that we are unable to process any new input // and we must wait for our in flight asynchronous operation to complete before // we are able to process more input. - ((InternalContext) context).awaitAsyncOperationComplete(); + context.awaitAsyncOperationComplete(); } } private void onAsyncResult( - Context context, AsyncOperationResult asyncResult) { + InternalContext context, AsyncOperationResult asyncResult) { if (asyncResult.unknown()) { ToFunction batch = asyncResult.metadata(); sendToFunction(context, batch); @@ -125,10 +127,10 @@ private void onAsyncResult( InvocationResponse invocationResult = unpackInvocationOrThrow(context.self(), asyncResult); handleInvocationResponse(context, invocationResult); - final int state = requestState.getOrDefault(-1); - if (state < 0) { + final int numBatched = requestState.getOrDefault(-1); + if (numBatched < 0) { throw new IllegalStateException("Got an unexpected async result"); - } else if (state == 0) { + } else if (numBatched == 0) { requestState.clear(); } else { final InvocationBatchRequest.Builder nextBatch = getNextBatch(); @@ -139,6 +141,7 @@ private void onAsyncResult( // b) sending the accumulated batch to the remote function. requestState.set(0); batch.clear(); + context.functionTypeMetrics().consumeBacklogMessages(numBatched); sendToFunction(context, nextBatch); } } diff --git a/statefun-flink/statefun-flink-core/src/test/java/org/apache/flink/statefun/flink/core/reqreply/RequestReplyFunctionTest.java b/statefun-flink/statefun-flink-core/src/test/java/org/apache/flink/statefun/flink/core/reqreply/RequestReplyFunctionTest.java index cd705888b..d4d341a53 100644 --- a/statefun-flink/statefun-flink-core/src/test/java/org/apache/flink/statefun/flink/core/reqreply/RequestReplyFunctionTest.java +++ b/statefun-flink/statefun-flink-core/src/test/java/org/apache/flink/statefun/flink/core/reqreply/RequestReplyFunctionTest.java @@ -209,6 +209,32 @@ public void egressIsSent() { new EgressIdentifier<>("org.foo", "bar", Any.class), context.egresses.get(0).getKey()); } + @Test + public void backlogMetricsIncreasedOnInvoke() { + functionUnderTest.invoke(context, Any.getDefaultInstance()); + + // following should be accounted into backlog metrics + functionUnderTest.invoke(context, Any.getDefaultInstance()); + functionUnderTest.invoke(context, Any.getDefaultInstance()); + + assertThat(context.functionTypeMetrics().numBacklog, is(2)); + } + + @Test + public void backlogMetricsDecreasedOnNextSuccess() { + functionUnderTest.invoke(context, Any.getDefaultInstance()); + + // following should be accounted into backlog metrics + functionUnderTest.invoke(context, Any.getDefaultInstance()); + functionUnderTest.invoke(context, Any.getDefaultInstance()); + + // complete one message, should fully consume backlog + context.needsWaiting = false; + functionUnderTest.invoke(context, successfulAsyncOperation()); + + assertThat(context.functionTypeMetrics().numBacklog, is(0)); + } + private static AsyncOperationResult successfulAsyncOperation() { return new AsyncOperationResult<>( new Object(), Status.SUCCESS, FromFunction.getDefaultInstance(), null); @@ -251,7 +277,7 @@ public int capturedInvocationBatchSize() { private static final class FakeContext implements InternalContext { - private final FunctionTypeMetrics fakeMetrics = new FakeMetrics(); + private final BacklogTrackingMetrics fakeMetrics = new BacklogTrackingMetrics(); Address caller; boolean needsWaiting; @@ -266,7 +292,7 @@ public void awaitAsyncOperationComplete() { } @Override - public FunctionTypeMetrics functionTypeMetrics() { + public BacklogTrackingMetrics functionTypeMetrics() { return fakeMetrics; } @@ -297,7 +323,23 @@ public void sendAfter(Duration delay, Address to, Object message) { public void registerAsyncOperation(M metadata, CompletableFuture future) {} } - private static final class FakeMetrics implements FunctionTypeMetrics { + private static final class BacklogTrackingMetrics implements FunctionTypeMetrics { + + private int numBacklog = 0; + + public int numBacklog() { + return numBacklog; + } + + @Override + public void appendBacklogMessages(int count) { + numBacklog += count; + } + + @Override + public void consumeBacklogMessages(int count) { + numBacklog -= count; + } @Override public void asyncOperationRegistered() {} @@ -322,11 +364,5 @@ public void blockedAddress() {} @Override public void unblockedAddress() {} - - @Override - public void appendBacklogMessages(int count) {} - - @Override - public void consumeBacklogMessages(int count) {} } }