|
27 | 27 | import org.apache.kafka.streams.processor.AbstractProcessor;
|
28 | 28 | import org.apache.kafka.streams.processor.Processor;
|
29 | 29 | import org.apache.kafka.streams.processor.ProcessorContext;
|
30 |
| -import org.apache.kafka.streams.processor.internals.InternalProcessorContext; |
31 | 30 | import org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl;
|
32 | 31 | import org.apache.kafka.streams.state.KeyValueIterator;
|
33 | 32 | import org.apache.kafka.streams.state.SessionStore;
|
|
38 | 37 | import java.util.ArrayList;
|
39 | 38 | import java.util.List;
|
40 | 39 |
|
41 |
| -import static org.apache.kafka.streams.processor.internals.metrics.TaskMetrics.droppedRecordsSensorOrLateRecordDropSensor; |
42 |
| -import static org.apache.kafka.streams.processor.internals.metrics.TaskMetrics.droppedRecordsSensorOrSkippedRecordsSensor; |
| 40 | +import static org.apache.kafka.streams.processor.internals.metrics.TaskMetrics.droppedRecordsSensor; |
43 | 41 |
|
44 | 42 | public class KStreamSessionWindowAggregate<K, V, Agg> implements KStreamAggProcessorSupplier<K, Windowed<K>, V, Agg> {
|
45 | 43 | private static final Logger LOG = LoggerFactory.getLogger(KStreamSessionWindowAggregate.class);
|
@@ -91,13 +89,12 @@ public void init(final ProcessorContext context) {
|
91 | 89 | super.init(context);
|
92 | 90 | final StreamsMetricsImpl metrics = (StreamsMetricsImpl) context.metrics();
|
93 | 91 | final String threadId = Thread.currentThread().getName();
|
94 |
| - lateRecordDropSensor = droppedRecordsSensorOrLateRecordDropSensor( |
| 92 | + lateRecordDropSensor = droppedRecordsSensor( |
95 | 93 | threadId,
|
96 | 94 | context.taskId().toString(),
|
97 |
| - ((InternalProcessorContext) context).currentNode().name(), |
98 | 95 | metrics
|
99 | 96 | );
|
100 |
| - droppedRecordsSensor = droppedRecordsSensorOrSkippedRecordsSensor(threadId, context.taskId().toString(), metrics); |
| 97 | + droppedRecordsSensor = droppedRecordsSensor(threadId, context.taskId().toString(), metrics); |
101 | 98 | store = context.getStateStore(storeName);
|
102 | 99 | tupleForwarder = new SessionTupleForwarder<>(store, context, new SessionCacheFlushListener<>(context), sendOldValues);
|
103 | 100 | }
|
|
0 commit comments