Refactoring with Java 8 Streams

Viewed 202

I implemented the following class DateBucket that has a method called bucketize. This method returns a list of DateBucket where each DateBucket in the list is between the dates fromDate and toDate and uses bucketSize and bucketSizeUnit as steps to set the values of the from and to DateBucket variables.

class DateBucket {
    final Instant from;
    final Instant to;

    public static List<DateBucket> bucketize(
        ZonedDateTime fromDate, 
        ZonedDateTime toDate, 
        int bucketSize, 
        ChronoUnit bucketSizeUnit
        ) {
    List<DateBucket> buckets = new ArrayList<>();
    
    boolean reachedDate = false;
    
    for (int i = 0; !reachedDate; i++) {
        ZonedDateTime minDate = fromDate.plus(i * bucketSize, bucketSizeUnit);
        ZonedDateTime maxDate = fromDate.plus((i + 1) * bucketSize, bucketSizeUnit);
        reachedDate = toDate.isBefore(maxDate);
        buckets.add(new DateBucket(minDate.toInstant(), maxDate.toInstant()));
    }
    return buckets;
}
}

As an example the following input:

fromDate -> 2020-12-07 00:00:00
toDate -> 2020-12-07 00:00:03
bucketSize -> 1
bucketSizeUnit -> ChronoUnit.SECONDS

Returns the following DateBucket list:

0. DateBucket(from = 2020-12-07  05:00:00, to = 2020-12-07 05:00:01)
1. DateBucket(from = 2020-12-07  05:00:01, to = 2020-12-07 05:00:02)
2. DateBucket(from = 2020-12-07  05:00:02, to = 2020-12-07 05:00:03)
3. DateBucket(from = 2020-12-07  05:00:03, to = 2020-12-07 05:00:04)

So, How can I refactor the bucketize method using the same logic but with Java 8 Streams?

3 Answers

You could do:

public static List<DateBucket> bucketize(
        ZonedDateTime fromDate, 
        ZonedDateTime toDate, 
        int bucketSize, 
        ChronoUnit bucketSizeUnit
    ) {
    
    var result = new ArrayList<DateBucket>();
    IntStream.range(0, Integer.MAX_VALUE)
            .mapToObj(i -> fromDate.plus(i * bucketSize, bucketSizeUnit))
            .takeWhile(d -> d.isBefore(toDate))
            .map(d -> d.toInstant())
            .reduce((min, max) -> {
                result.add(new DateBucket(min, max));
                return max;
            });
    return result;
}

While this reduction function could be considered an abuse of the stream API (reduction functions are supposed to be stateless and associative, which ours is not, but it should be safe in this case since we know that the stream is sequential, and will thus iterate in order), the stream API doesn't seem to offer a more straightforward API to combine adjacent stream elements.

That said, I'd probably prefer your imperative implementation.

The line of thought here is that you want to stream the values needed to create the buckets, then map and collect them. A stream of ZonedDateTimes will do and you can use a Spliterator to create the stream. This seems to be one possible implementation which is basic and doesn't support parallelization, but does seem to work.

public class ZonedDateTimeSpliterator implements Spliterator<ZonedDateTime> {
    private final ZonedDateTime fromDate;
    private final ZonedDateTime toDate;
    private final int bucketSize;
    private final ChronoUnit bucketSizeUnit;
    private ZonedDateTime currentDate;

    public ZonedDateTimeSpliterator(ZonedDateTime fromDate, ZonedDateTime toDate, int bucketSize,
            ChronoUnit bucketSizeUnit) {
        this.fromDate = fromDate;
        this.toDate = toDate;
        this.bucketSize = bucketSize;
        this.bucketSizeUnit = bucketSizeUnit;
        currentDate = fromDate.truncatedTo(bucketSizeUnit);
    }

    public void forEachRemaining(Consumer<? super ZonedDateTime> action) {
        while (currentDate.plus(bucketSize, bucketSizeUnit).isBefore(toDate)) {
            getNext(action);
        }
    }

    private void getNext(Consumer<? super ZonedDateTime> action) {
        action.accept(currentDate);
        currentDate = currentDate.plus(bucketSize, bucketSizeUnit);
    }

    @Override
    public boolean tryAdvance(Consumer<? super ZonedDateTime> action) {
        if (currentDate.plus(bucketSize, bucketSizeUnit).isBefore(toDate)) {
            getNext(action);
            return true;
        } else // cannot advance
            return false;
    }

    @Override
    public Spliterator<ZonedDateTime> trySplit() {
        // implement for parallel processing
        return null;
    }

    @Override
    public long estimateSize() {
        return Duration.between(fromDate.truncatedTo(bucketSizeUnit), toDate.truncatedTo(bucketSizeUnit))
                .get(bucketSizeUnit);
    }

    @Override
    public int characteristics() {
        return ORDERED | SIZED | IMMUTABLE;
    }
}

Then you simply stream the values, map the buckets, and collect the list.

private List<DateBucket> doThing(ZonedDateTime fromDate, ZonedDateTime toDate, int bucketSize,
        ChronoUnit bucketSizeUnit) {
    return StreamSupport
            .stream(() -> new ZonedDateTimeSpliterator(fromDate, toDate, bucketSize, bucketSizeUnit),
                    Spliterator.ORDERED | Spliterator.SIZED | Spliterator.IMMUTABLE, false)
            .map(currentDate -> new DateBucket(currentDate.toInstant(),
                    currentDate.plus(bucketSize, bucketSizeUnit).toInstant()))
            .collect(Collectors.toList());

}

You could range over the number of required buckets by calculating this number in advance using TemporalUnit#between:

public static List<DateBucket> bucketize(ZonedDateTime fromDate, ZonedDateTime toDate,
    int bucketSize, ChronoUnit bucketSizeUnit) {

    return LongStream.range(0, bucketSizeUnit.between(fromDate, toDate)/bucketSize + 1)
        .mapToObj( nb -> new DateBucket(
            fromDate.plus(nb * bucketSize, bucketSizeUnit).toInstant(),
            fromDate.plus((nb+1) * bucketSize, bucketSizeUnit).toInstant()))
        .collect(Collectors.toList());
}
Related