Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -104,8 +104,8 @@ Constructor <org.apache.flink.connector.file.src.AbstractFileSource$AbstractFile
Constructor <org.apache.flink.connector.file.src.AbstractFileSource.<init>([Lorg.apache.flink.core.fs.Path;, org.apache.flink.connector.file.src.enumerate.FileEnumerator$Provider, org.apache.flink.connector.file.src.assigners.FileSplitAssigner$Provider, org.apache.flink.connector.file.src.reader.BulkFormat, org.apache.flink.connector.file.src.ContinuousEnumerationSettings)> has parameter of type <[Lorg.apache.flink.core.fs.Path;> in (AbstractFileSource.java:0)
Constructor <org.apache.flink.connector.file.src.FileSource$FileSourceBuilder.<init>([Lorg.apache.flink.core.fs.Path;, org.apache.flink.connector.file.src.reader.BulkFormat)> has parameter of type <[Lorg.apache.flink.core.fs.Path;> in (FileSource.java:0)
Constructor <org.apache.flink.connector.file.src.FileSource.<init>([Lorg.apache.flink.core.fs.Path;, org.apache.flink.connector.file.src.enumerate.FileEnumerator$Provider, org.apache.flink.connector.file.src.assigners.FileSplitAssigner$Provider, org.apache.flink.connector.file.src.reader.BulkFormat, org.apache.flink.connector.file.src.ContinuousEnumerationSettings)> has parameter of type <[Lorg.apache.flink.core.fs.Path;> in (FileSource.java:0)
Constructor <org.apache.flink.connector.file.src.assigners.LocalityAwareSplitAssigner.<init>(java.util.Collection)> calls constructor <org.apache.flink.metrics.SimpleCounter.<init>()> in (LocalityAwareSplitAssigner.java:80)
Constructor <org.apache.flink.connector.file.src.assigners.LocalityAwareSplitAssigner.<init>(java.util.Collection)> calls constructor <org.apache.flink.metrics.SimpleCounter.<init>()> in (LocalityAwareSplitAssigner.java:81)
Constructor <org.apache.flink.connector.file.src.assigners.LocalityAwareSplitAssigner.<init>(java.util.Collection)> calls constructor <org.apache.flink.metrics.SimpleMonotonicCounter.<init>()> in (LocalityAwareSplitAssigner.java:80)
Constructor <org.apache.flink.connector.file.src.assigners.LocalityAwareSplitAssigner.<init>(java.util.Collection)> calls constructor <org.apache.flink.metrics.SimpleMonotonicCounter.<init>()> in (LocalityAwareSplitAssigner.java:81)
Constructor <org.apache.flink.connector.file.src.impl.ContinuousFileSplitEnumerator.<init>(org.apache.flink.api.connector.source.SplitEnumeratorContext, org.apache.flink.connector.file.src.enumerate.FileEnumerator, org.apache.flink.connector.file.src.assigners.FileSplitAssigner, [Lorg.apache.flink.core.fs.Path;, java.util.Collection, long)> has parameter of type <[Lorg.apache.flink.core.fs.Path;> in (ContinuousFileSplitEnumerator.java:0)
Constructor <org.apache.flink.connector.file.table.ColumnarRowIterator.<init>(org.apache.flink.table.data.columnar.ColumnarRowData, java.lang.Runnable)> has parameter of type <org.apache.flink.table.data.columnar.ColumnarRowData> in (ColumnarRowIterator.java:0)
Constructor <org.apache.flink.connector.file.table.FileSystemOutputFormat$Builder.<init>()> calls constructor <org.apache.flink.streaming.api.functions.sink.filesystem.OutputFileConfig.<init>(java.lang.String, java.lang.String)> in (FileSystemOutputFormat.java:237)
Expand Down Expand Up @@ -163,8 +163,8 @@ Field <org.apache.flink.connector.file.sink.writer.FileWriterBucketStateSerializ
Field <org.apache.flink.connector.file.src.AbstractFileSource$AbstractFileSourceBuilder.inputPaths> has type <[Lorg.apache.flink.core.fs.Path;> in (AbstractFileSource.java:0)
Field <org.apache.flink.connector.file.src.AbstractFileSource.inputPaths> has type <[Lorg.apache.flink.core.fs.Path;> in (AbstractFileSource.java:0)
Field <org.apache.flink.connector.file.src.FileSourceSplitSerializer.SERIALIZER_CACHE> has generic type <java.lang.ThreadLocal<org.apache.flink.core.memory.DataOutputSerializer>> with type argument depending on <org.apache.flink.core.memory.DataOutputSerializer> in (FileSourceSplitSerializer.java:0)
Field <org.apache.flink.connector.file.src.assigners.LocalityAwareSplitAssigner.localAssignments> has type <org.apache.flink.metrics.SimpleCounter> in (LocalityAwareSplitAssigner.java:0)
Field <org.apache.flink.connector.file.src.assigners.LocalityAwareSplitAssigner.remoteAssignments> has type <org.apache.flink.metrics.SimpleCounter> in (LocalityAwareSplitAssigner.java:0)
Field <org.apache.flink.connector.file.src.assigners.LocalityAwareSplitAssigner.localAssignments> has type <org.apache.flink.metrics.SimpleMonotonicCounter> in (LocalityAwareSplitAssigner.java:0)
Field <org.apache.flink.connector.file.src.assigners.LocalityAwareSplitAssigner.remoteAssignments> has type <org.apache.flink.metrics.SimpleMonotonicCounter> in (LocalityAwareSplitAssigner.java:0)
Field <org.apache.flink.connector.file.src.compression.StandardDeCompressors.DECOMPRESSORS> has generic type <java.util.Map<java.lang.String, org.apache.flink.api.common.io.compression.InflaterInputStreamFactory<?>>> with type argument depending on <org.apache.flink.api.common.io.compression.InflaterInputStreamFactory> in (StandardDeCompressors.java:0)
Field <org.apache.flink.connector.file.src.impl.ContinuousFileSplitEnumerator.paths> has type <[Lorg.apache.flink.core.fs.Path;> in (ContinuousFileSplitEnumerator.java:0)
Field <org.apache.flink.connector.file.table.ColumnarRowIterator.rowData> has type <org.apache.flink.table.data.columnar.ColumnarRowData> in (ColumnarRowIterator.java:0)
Expand Down Expand Up @@ -492,12 +492,12 @@ Method <org.apache.flink.connector.file.src.FileSourceSplitSerializer.serialize(
Method <org.apache.flink.connector.file.src.FileSourceSplitSerializer.serialize(org.apache.flink.connector.file.src.FileSourceSplit)> calls method <org.apache.flink.core.memory.DataOutputSerializer.writeLong(long)> in (FileSourceSplitSerializer.java:77)
Method <org.apache.flink.connector.file.src.FileSourceSplitSerializer.serialize(org.apache.flink.connector.file.src.FileSourceSplit)> calls method <org.apache.flink.core.memory.DataOutputSerializer.writeLong(long)> in (FileSourceSplitSerializer.java:78)
Method <org.apache.flink.connector.file.src.FileSourceSplitSerializer.serialize(org.apache.flink.connector.file.src.FileSourceSplit)> calls method <org.apache.flink.core.memory.DataOutputSerializer.writeUTF(java.lang.String)> in (FileSourceSplitSerializer.java:66)
Method <org.apache.flink.connector.file.src.assigners.LocalityAwareSplitAssigner.getNext(java.lang.String)> calls method <org.apache.flink.metrics.SimpleCounter.inc()> in (LocalityAwareSplitAssigner.java:114)
Method <org.apache.flink.connector.file.src.assigners.LocalityAwareSplitAssigner.getNumberOfLocalAssignments()> calls method <org.apache.flink.metrics.SimpleCounter.getCount()> in (LocalityAwareSplitAssigner.java:156)
Method <org.apache.flink.connector.file.src.assigners.LocalityAwareSplitAssigner.getNext(java.lang.String)> calls method <org.apache.flink.metrics.SimpleMonotonicCounter.inc()> in (LocalityAwareSplitAssigner.java:114)
Method <org.apache.flink.connector.file.src.assigners.LocalityAwareSplitAssigner.getNumberOfLocalAssignments()> calls method <org.apache.flink.metrics.SimpleMonotonicCounter.getCount()> in (LocalityAwareSplitAssigner.java:156)
Method <org.apache.flink.connector.file.src.assigners.LocalityAwareSplitAssigner.getNumberOfLocalAssignments()> is annotated with <org.apache.flink.annotation.VisibleForTesting> in (LocalityAwareSplitAssigner.java:0)
Method <org.apache.flink.connector.file.src.assigners.LocalityAwareSplitAssigner.getNumberOfRemoteAssignments()> calls method <org.apache.flink.metrics.SimpleCounter.getCount()> in (LocalityAwareSplitAssigner.java:161)
Method <org.apache.flink.connector.file.src.assigners.LocalityAwareSplitAssigner.getNumberOfRemoteAssignments()> calls method <org.apache.flink.metrics.SimpleMonotonicCounter.getCount()> in (LocalityAwareSplitAssigner.java:161)
Method <org.apache.flink.connector.file.src.assigners.LocalityAwareSplitAssigner.getNumberOfRemoteAssignments()> is annotated with <org.apache.flink.annotation.VisibleForTesting> in (LocalityAwareSplitAssigner.java:0)
Method <org.apache.flink.connector.file.src.assigners.LocalityAwareSplitAssigner.getRemoteSplit()> calls method <org.apache.flink.metrics.SimpleCounter.inc()> in (LocalityAwareSplitAssigner.java:150)
Method <org.apache.flink.connector.file.src.assigners.LocalityAwareSplitAssigner.getRemoteSplit()> calls method <org.apache.flink.metrics.SimpleMonotonicCounter.inc()> in (LocalityAwareSplitAssigner.java:150)
Method <org.apache.flink.connector.file.src.compression.StandardDeCompressors.buildDecompressorMap([Lorg.apache.flink.api.common.io.compression.InflaterInputStreamFactory;)> calls method <org.apache.flink.api.common.io.compression.InflaterInputStreamFactory.getCommonFileExtensions()> in (StandardDeCompressors.java:91)
Method <org.apache.flink.connector.file.src.compression.StandardDeCompressors.buildDecompressorMap([Lorg.apache.flink.api.common.io.compression.InflaterInputStreamFactory;)> depends on component type <org.apache.flink.api.common.io.compression.InflaterInputStreamFactory> in (StandardDeCompressors.java:0)
Method <org.apache.flink.connector.file.src.compression.StandardDeCompressors.buildDecompressorMap([Lorg.apache.flink.api.common.io.compression.InflaterInputStreamFactory;)> has generic return type <java.util.Map<java.lang.String, org.apache.flink.api.common.io.compression.InflaterInputStreamFactory<?>>> with type argument depending on <org.apache.flink.api.common.io.compression.InflaterInputStreamFactory> in (StandardDeCompressors.java:0)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@
import org.apache.flink.annotation.PublicEvolving;
import org.apache.flink.annotation.VisibleForTesting;
import org.apache.flink.connector.file.src.FileSourceSplit;
import org.apache.flink.metrics.SimpleCounter;
import org.apache.flink.metrics.SimpleMonotonicCounter;
import org.apache.flink.util.MathUtils;
import org.apache.flink.util.NetUtils;
import org.apache.flink.util.StringUtils;
Expand Down Expand Up @@ -64,8 +64,8 @@ public class LocalityAwareSplitAssigner implements FileSplitAssigner {
/** Unassigned splits for remote assignment. */
private final LocatableSplitChooser remoteSplitChooser;

private final SimpleCounter localAssignments;
private final SimpleCounter remoteAssignments;
private final SimpleMonotonicCounter localAssignments;
private final SimpleMonotonicCounter remoteAssignments;

// --------------------------------------------------------------------------------------------

Expand All @@ -77,8 +77,8 @@ public LocalityAwareSplitAssigner(Collection<FileSourceSplit> splits) {

// this will be replaced with metrics registration once we can expose the metric group
// properly to the assigners
this.localAssignments = new SimpleCounter();
this.remoteAssignments = new SimpleCounter();
this.localAssignments = new SimpleMonotonicCounter();
this.remoteAssignments = new SimpleMonotonicCounter();
}

// --------------------------------------------------------------------------------------------
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@
import org.apache.flink.metrics.Gauge;
import org.apache.flink.metrics.Histogram;
import org.apache.flink.metrics.MetricGroup;
import org.apache.flink.metrics.ThreadSafeSimpleCounter;
import org.apache.flink.metrics.ThreadSafeSimpleMonotonicCounter;
import org.apache.flink.runtime.metrics.DescriptiveStatisticsHistogram;
import org.apache.flink.runtime.metrics.groups.ProxyMetricGroup;

Expand All @@ -46,7 +46,9 @@ public class ChangelogStorageMetricGroup extends ProxyMetricGroup<MetricGroup> {
public ChangelogStorageMetricGroup(MetricGroup parent) {
super(parent);
this.uploadsCounter =
counter(CHANGELOG_STORAGE_NUM_UPLOAD_REQUESTS, new ThreadSafeSimpleCounter());
counter(
CHANGELOG_STORAGE_NUM_UPLOAD_REQUESTS,
new ThreadSafeSimpleMonotonicCounter());
this.uploadBatchSizes =
histogram(
CHANGELOG_STORAGE_UPLOAD_BATCH_SIZES,
Expand All @@ -68,7 +70,9 @@ public ChangelogStorageMetricGroup(MetricGroup parent) {
CHANGELOG_STORAGE_UPLOAD_LATENCIES_NANOS,
new DescriptiveStatisticsHistogram(WINDOW_SIZE));
this.uploadFailuresCounter =
counter(CHANGELOG_STORAGE_NUM_UPLOAD_FAILURES, new ThreadSafeSimpleCounter());
counter(
CHANGELOG_STORAGE_NUM_UPLOAD_FAILURES,
new ThreadSafeSimpleMonotonicCounter());
}

public Counter getUploadsCounter() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -34,14 +34,33 @@ public interface Counter extends Metric {
*/
void inc(long n);

/** Decrement the current count by 1. */
/**
* Decrement the current count by 1.
*
* <p>This is an <i>optional operation</i>: implementations whose count is meant to only ever
* increase (see {@link MonotonicCounter}) may refuse to support it and throw {@link
* UnsupportedOperationException} instead.
*
* @throws UnsupportedOperationException if this counter does not support decrementing
* @deprecated If you need a counter that can decrement, please migrate to {@link
* UpDownCounter}, where decrementing is guaranteed to be fully supported.
*/
@Deprecated
void dec();

/**
* Decrement the current count by the given value.
*
* <p>This is an <i>optional operation</i>: implementations whose count is meant to only ever
* increase (see {@link MonotonicCounter}) may refuse to support it and throw {@link
* UnsupportedOperationException} instead.
*
* @param n value to decrement the current count by
* @throws UnsupportedOperationException if this counter does not support decrementing
* @deprecated If you need a counter that can decrement, please migrate to {@link
* UpDownCounter}, where decrementing is guaranteed to be fully supported.
*/
@Deprecated
void dec(long n);

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@

import org.apache.flink.annotation.Experimental;
import org.apache.flink.annotation.Public;
import org.apache.flink.annotation.PublicEvolving;
import org.apache.flink.events.EventBuilder;
import org.apache.flink.traces.SpanBuilder;

Expand Down Expand Up @@ -94,6 +95,32 @@ default <C extends Counter> C counter(int name, C counter) {
*/
<C extends Counter> C counter(String name, C counter);

/**
* Creates and registers a new {@link MonotonicCounter} with Flink.
*
* <p>Use this for counts that only ever increase, so that reporters can export them as
* monotonic/cumulative counters and benefit from automatic reset detection.
*
* @param name name of the counter
* @return the created counter
*/
@PublicEvolving
default MonotonicCounter monotonicCounter(String name) {
return counter(name, new SimpleMonotonicCounter());
}

/**
* Creates and registers a new {@link MonotonicCounter} with Flink.
*
* @param name name of the counter
* @return the created counter
* @see #monotonicCounter(String)
*/
@PublicEvolving
default MonotonicCounter monotonicCounter(int name) {
return monotonicCounter(String.valueOf(name));
}

/**
* Registers a new {@link org.apache.flink.metrics.Gauge} with Flink.
*
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,61 @@
/*
* 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.metrics;

import org.apache.flink.annotation.PublicEvolving;

/**
* A {@link Counter} whose count only ever increases; it is never decremented.
*
* <p>Reporters that want to export this as a monotonic/cumulative counter in the target monitoring
* system (e.g. an OpenTelemetry monotonic {@code Sum} or a native Prometheus counter, enabling
* automatic reset detection) should check for this marker interface with {@code instanceof} rather
* than relying on {@link Metric#getMetricType()}, which continues to report {@link
* MetricType#COUNTER} for a {@code MonotonicCounter} so that reporters unaware of this interface
* keep treating it as a regular counter.
*
* <p>Decrementing is one of {@link Counter}'s optional operations (see {@link Counter#dec()}):
* {@link #dec()} and {@link #dec(long)} throw {@link UnsupportedOperationException} rather than
* silently violating monotonicity.
*/
@PublicEvolving
public interface MonotonicCounter extends Counter {

/**
* Always throws {@link UnsupportedOperationException}, since a {@code MonotonicCounter}'s count
* only ever increases.
*
* @throws UnsupportedOperationException always
*/
@Override
default void dec() {
throw new UnsupportedOperationException("MonotonicCounter does not support decrementing.");
}

/**
* Always throws {@link UnsupportedOperationException}, since a {@code MonotonicCounter}'s count
* only ever increases.
*
* @throws UnsupportedOperationException always
*/
@Override
default void dec(long n) {
throw new UnsupportedOperationException("MonotonicCounter does not support decrementing.");
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -20,9 +20,9 @@

import org.apache.flink.annotation.Internal;

/** A simple low-overhead {@link org.apache.flink.metrics.Counter} that is not thread-safe. */
/** A simple low-overhead {@link org.apache.flink.metrics.UpDownCounter} that is not thread-safe. */
@Internal
public class SimpleCounter implements Counter {
public class SimpleCounter implements UpDownCounter {

/** the current count. */
private long count;
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,61 @@
/*
* 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.metrics;

import org.apache.flink.annotation.Internal;

/**
* A simple low-overhead {@link MonotonicCounter} that is not thread-safe.
*
* <p>It behaves like {@link SimpleCounter} for {@link #inc()} / {@link #inc(long)} / {@link
* #getCount()}, and rejects {@link #dec()} / {@link #dec(long)} per the {@link MonotonicCounter}
* contract.
*/
@Internal
public class SimpleMonotonicCounter implements MonotonicCounter {

/** the current count. */
private long count;

/** Increment the current count by 1. */
@Override
public void inc() {
count++;
}

/**
* Increment the current count by the given value.
*
* @param n value to increment the current count by
*/
@Override
public void inc(long n) {
count += n;
}

/**
* Returns the current count.
*
* @return current count
*/
@Override
public long getCount() {
return count;
}
}
Loading