Skip to content
This repository has been archived by the owner on Mar 15, 2021. It is now read-only.

Commit

Permalink
Handling high freq in batches (#253)
Browse files Browse the repository at this point in the history
* Handling high freq in batches

* Fix stats reporting
  • Loading branch information
Sergey Rustamov authored Mar 2, 2021
1 parent fb7922d commit 61a305b
Showing 1 changed file with 9 additions and 10 deletions.
19 changes: 9 additions & 10 deletions api/src/main/java/com/spotify/ffwd/output/BatchingPluginSink.java
Original file line number Diff line number Diff line change
Expand Up @@ -23,21 +23,19 @@
import com.google.inject.Inject;
import com.google.inject.name.Named;
import com.spotify.ffwd.filter.Filter;
import com.spotify.ffwd.model.v2.Batch;
import com.spotify.ffwd.model.v2.Metric;
import com.spotify.ffwd.statistics.BatchingStatistics;
import com.spotify.ffwd.statistics.OutputPluginStatistics;
import com.spotify.ffwd.util.BatchMetricConverter;
import com.spotify.ffwd.util.HighFrequencyDetector;
import eu.toolchain.async.AsyncFramework;
import eu.toolchain.async.AsyncFuture;
import eu.toolchain.async.FutureFinished;
import eu.toolchain.async.LazyTransform;
import eu.toolchain.async.Transform;
import java.util.ArrayList;
import java.util.Date;
import java.util.HashSet;
import java.util.List;
import java.util.Optional;
import java.util.Set;

import java.util.*;
import java.util.concurrent.Callable;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.ScheduledFuture;
Expand Down Expand Up @@ -384,12 +382,13 @@ AsyncFuture<Void> doFlush(Batch newBatch) {
.onFinished(() -> batchingStatistics.reportSentMetrics(batch.metrics.size())));
}

// TODO finish batches
if (!batch.batches.isEmpty()) {
final List<Metric> metrics = BatchMetricConverter.convertBatchesToMetrics(batch.batches);

final List<Metric> filteredMetrics = highFrequencyDetector.detect(metrics);
futures.add(sink
.sendBatches(batch.batches)
.onFinished(() -> batchingStatistics.reportSentBatches(batch.batches.size(),
batch.size())));
.sendMetrics(filteredMetrics)
.onFinished(() -> batchingStatistics.reportSentMetrics(filteredMetrics.size())));
}

// chain into batch future.
Expand Down

0 comments on commit 61a305b

Please sign in to comment.