|
1 | 1 | package org.opensearch.dataprepper.plugins.source.source_crawler.base; |
2 | 2 |
|
3 | 3 | import io.micrometer.core.instrument.Counter; |
| 4 | +import io.micrometer.core.instrument.Timer; |
4 | 5 | import org.opensearch.dataprepper.metrics.PluginMetrics; |
5 | 6 | import org.opensearch.dataprepper.model.acknowledgements.AcknowledgementSet; |
6 | 7 | import org.opensearch.dataprepper.model.buffer.Buffer; |
|
36 | 37 | public class DimensionalTimeSliceCrawler implements Crawler<DimensionalTimeSliceWorkerProgressState> { |
37 | 38 | private static final Logger log = LoggerFactory.getLogger(DimensionalTimeSliceCrawler.class); |
38 | 39 | private static final String DIMENSIONAL_TIME_SLICE_WORKER_PARTITIONS_CREATED = "DimensionalTimeSliceWorkerPartitionsCreated"; |
| 40 | + private static final String WORKER_PARTITION_WAIT_TIME = "WorkerPartitionWaitTime"; |
| 41 | + private static final String WORKER_PARTITION_PROCESS_LATENCY = "WorkerPartitionProcessLatency"; |
39 | 42 | private static final Duration HOUR_DURATION = Duration.ofHours(1); |
40 | 43 |
|
41 | 44 | private final CrawlerClient client; |
42 | 45 | private final Counter partitionsCreatedCounter; |
| 46 | + private final Timer partitionWaitTimeTimer; |
| 47 | + private final Timer partitionProcessLatencyTimer; |
43 | 48 | private List<String> dimensionTypes; |
44 | 49 | private static final String LAST_UPDATED_KEY = "last_updated|"; |
45 | 50 |
|
46 | 51 | public DimensionalTimeSliceCrawler(CrawlerClient client, |
47 | 52 | PluginMetrics pluginMetrics) { |
48 | 53 | this.client = client; |
49 | 54 | this.partitionsCreatedCounter = pluginMetrics.counter(DIMENSIONAL_TIME_SLICE_WORKER_PARTITIONS_CREATED); |
| 55 | + this.partitionWaitTimeTimer = pluginMetrics.timer(WORKER_PARTITION_WAIT_TIME); |
| 56 | + this.partitionProcessLatencyTimer = pluginMetrics.timer(WORKER_PARTITION_PROCESS_LATENCY); |
50 | 57 | } |
51 | 58 |
|
52 | 59 | /** |
@@ -78,7 +85,8 @@ public Instant crawl(LeaderPartition leaderPartition, EnhancedSourceCoordinator |
78 | 85 |
|
79 | 86 | @Override |
80 | 87 | public void executePartition(DimensionalTimeSliceWorkerProgressState state, Buffer<Record<Event>> buffer, AcknowledgementSet acknowledgementSet) { |
81 | | - client.executePartition(state, buffer, acknowledgementSet); |
| 88 | + partitionWaitTimeTimer.record(Duration.between(state.getPartitionCreationTime(), Instant.now())); |
| 89 | + partitionProcessLatencyTimer.record(() -> client.executePartition(state, buffer, acknowledgementSet)); |
82 | 90 | } |
83 | 91 |
|
84 | 92 | private void createPartitionsForDimensionTypes(LeaderPartition leaderPartition, |
@@ -149,6 +157,7 @@ private void createPartitionForIncrementalSync(LeaderPartition leaderPartition, |
149 | 157 | void createWorkerPartition(Instant startTime, Instant endTime, |
150 | 158 | String dimensionType, EnhancedSourceCoordinator coordinator) { |
151 | 159 | DimensionalTimeSliceWorkerProgressState workerState = new DimensionalTimeSliceWorkerProgressState(); |
| 160 | + workerState.setPartitionCreationTime(Instant.now()); |
152 | 161 | workerState.setStartTime(startTime); |
153 | 162 | workerState.setEndTime(endTime); |
154 | 163 | workerState.setDimensionType(dimensionType); |
|
0 commit comments