!1015 【6.0】处理EasyStream并行流下zip等依赖于toIdxMap函数时没有保证顺序问题导致的数据顺序错乱的问题

Merge pull request !1015 from 阿超/v6-dev
This commit is contained in:
Looly 2023-06-13 03:48:40 +00:00 committed by Gitee
commit 52feb1b60f
No known key found for this signature in database
GPG Key ID: 173E9B9CA92EEF8F
3 changed files with 28 additions and 94 deletions

View File

@ -12,10 +12,10 @@
package org.dromara.hutool.core.stream; package org.dromara.hutool.core.stream;
import org.dromara.hutool.core.array.ArrayUtil;
import org.dromara.hutool.core.lang.Opt; import org.dromara.hutool.core.lang.Opt;
import org.dromara.hutool.core.lang.mutable.MutableInt; import org.dromara.hutool.core.lang.mutable.MutableInt;
import org.dromara.hutool.core.lang.mutable.MutableObj; import org.dromara.hutool.core.lang.mutable.MutableObj;
import org.dromara.hutool.core.array.ArrayUtil;
import java.util.*; import java.util.*;
import java.util.function.*; import java.util.function.*;
@ -214,7 +214,7 @@ public interface TerminableWrappedStream<T, S extends TerminableWrappedStream<T,
*/ */
default <U> Map<Integer, U> toIdxMap(final Function<? super T, ? extends U> valueMapper) { default <U> Map<Integer, U> toIdxMap(final Function<? super T, ? extends U> valueMapper) {
final MutableInt index = new MutableInt(NOT_FOUND_ELEMENT_INDEX); final MutableInt index = new MutableInt(NOT_FOUND_ELEMENT_INDEX);
return EasyStream.of(toList()).toMap(e -> index.incrementAndGet(), valueMapper, (l, r) -> r); return EasyStream.of(sequential().toList()).toMap(e -> index.incrementAndGet(), valueMapper, (l, r) -> r);
} }
// region ============ to zip ============ // region ============ to zip ============
@ -269,16 +269,12 @@ public interface TerminableWrappedStream<T, S extends TerminableWrappedStream<T,
@SuppressWarnings("ResultOfMethodCallIgnored") @SuppressWarnings("ResultOfMethodCallIgnored")
default int findFirstIdx(final Predicate<? super T> predicate) { default int findFirstIdx(final Predicate<? super T> predicate) {
Objects.requireNonNull(predicate); Objects.requireNonNull(predicate);
if (isParallel()) { final MutableInt index = new MutableInt(NOT_FOUND_ELEMENT_INDEX);
return NOT_FOUND_ELEMENT_INDEX; unwrap().filter(e -> {
} else { index.increment();
final MutableInt index = new MutableInt(NOT_FOUND_ELEMENT_INDEX); return predicate.test(e);
unwrap().filter(e -> { }).findFirst();// 此处只做计数不需要值
index.increment(); return index.get();
return predicate.test(e);
}).findFirst();// 此处只做计数不需要值
return index.get();
}
} }
/** /**
@ -317,17 +313,13 @@ public interface TerminableWrappedStream<T, S extends TerminableWrappedStream<T,
*/ */
default int findLastIdx(final Predicate<? super T> predicate) { default int findLastIdx(final Predicate<? super T> predicate) {
Objects.requireNonNull(predicate); Objects.requireNonNull(predicate);
if (isParallel()) { final MutableInt idxRef = new MutableInt(NOT_FOUND_ELEMENT_INDEX);
return NOT_FOUND_ELEMENT_INDEX; forEachIdx((e, i) -> {
} else { if (predicate.test(e)) {
final MutableInt idxRef = new MutableInt(NOT_FOUND_ELEMENT_INDEX); idxRef.set(i);
forEachIdx((e, i) -> { }
if (predicate.test(e)) { });
idxRef.set(i); return idxRef.get();
}
});
return idxRef.get();
}
} }
/** /**
@ -505,13 +497,8 @@ public interface TerminableWrappedStream<T, S extends TerminableWrappedStream<T,
*/ */
default void forEachIdx(final BiConsumer<? super T, Integer> action) { default void forEachIdx(final BiConsumer<? super T, Integer> action) {
Objects.requireNonNull(action); Objects.requireNonNull(action);
if (isParallel()) { final MutableInt index = new MutableInt(NOT_FOUND_ELEMENT_INDEX);
EasyStream.of(toIdxMap().entrySet()).parallel(isParallel()) unwrap().forEach(e -> action.accept(e, index.incrementAndGet()));
.forEach(e -> action.accept(e.getValue(), e.getKey()));
} else {
final MutableInt index = new MutableInt(NOT_FOUND_ELEMENT_INDEX);
unwrap().forEach(e -> action.accept(e, index.incrementAndGet()));
}
} }
/** /**
@ -522,14 +509,8 @@ public interface TerminableWrappedStream<T, S extends TerminableWrappedStream<T,
*/ */
default void forEachOrderedIdx(final BiConsumer<? super T, Integer> action) { default void forEachOrderedIdx(final BiConsumer<? super T, Integer> action) {
Objects.requireNonNull(action); Objects.requireNonNull(action);
if (isParallel()) { final MutableInt index = new MutableInt(NOT_FOUND_ELEMENT_INDEX);
EasyStream.of(toIdxMap().entrySet()) unwrap().forEachOrdered(e -> action.accept(e, index.incrementAndGet()));
.parallel(isParallel())
.forEachOrdered(e -> action.accept(e.getValue(), e.getKey()));
} else {
final MutableInt index = new MutableInt(NOT_FOUND_ELEMENT_INDEX);
unwrap().forEachOrdered(e -> action.accept(e, index.incrementAndGet()));
}
} }
// endregion // endregion

View File

@ -12,6 +12,7 @@
package org.dromara.hutool.core.stream; package org.dromara.hutool.core.stream;
import org.dromara.hutool.core.array.ArrayUtil;
import org.dromara.hutool.core.collection.ListUtil; import org.dromara.hutool.core.collection.ListUtil;
import org.dromara.hutool.core.collection.iter.IterUtil; import org.dromara.hutool.core.collection.iter.IterUtil;
import org.dromara.hutool.core.lang.Console; import org.dromara.hutool.core.lang.Console;
@ -19,7 +20,6 @@ import org.dromara.hutool.core.lang.mutable.MutableInt;
import org.dromara.hutool.core.lang.mutable.MutableObj; import org.dromara.hutool.core.lang.mutable.MutableObj;
import org.dromara.hutool.core.map.MapUtil; import org.dromara.hutool.core.map.MapUtil;
import org.dromara.hutool.core.map.SafeConcurrentHashMap; import org.dromara.hutool.core.map.SafeConcurrentHashMap;
import org.dromara.hutool.core.array.ArrayUtil;
import java.util.*; import java.util.*;
import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicBoolean;
@ -279,16 +279,8 @@ public interface TransformableWrappedStream<T, S extends TransformableWrappedStr
*/ */
default S peekIdx(final BiConsumer<? super T, Integer> action) { default S peekIdx(final BiConsumer<? super T, Integer> action) {
Objects.requireNonNull(action); Objects.requireNonNull(action);
if (isParallel()) { final AtomicInteger index = new AtomicInteger(NOT_FOUND_ELEMENT_INDEX);
final Map<Integer, T> idxMap = easyStream().toIdxMap(); return peek(e -> action.accept(e, index.incrementAndGet()));
return wrap(EasyStream.of(idxMap.entrySet())
.parallel(isParallel())
.peek(e -> action.accept(e.getValue(), e.getKey()))
.map(Map.Entry::getValue));
} else {
final AtomicInteger index = new AtomicInteger(NOT_FOUND_ELEMENT_INDEX);
return peek(e -> action.accept(e, index.incrementAndGet()));
}
} }
/** /**
@ -383,16 +375,8 @@ public interface TransformableWrappedStream<T, S extends TransformableWrappedStr
*/ */
default S filterIdx(final BiPredicate<? super T, Integer> predicate) { default S filterIdx(final BiPredicate<? super T, Integer> predicate) {
Objects.requireNonNull(predicate); Objects.requireNonNull(predicate);
if (isParallel()) { final MutableInt index = new MutableInt(NOT_FOUND_ELEMENT_INDEX);
final Map<Integer, T> idxMap = easyStream().toIdxMap(); return filter(e -> predicate.test(e, index.incrementAndGet()));
return wrap(EasyStream.of(idxMap.entrySet())
.parallel(isParallel())
.filter(e -> predicate.test(e.getValue(), e.getKey()))
.map(Map.Entry::getValue));
} else {
final MutableInt index = new MutableInt(NOT_FOUND_ELEMENT_INDEX);
return filter(e -> predicate.test(e, index.incrementAndGet()));
}
} }
/** /**
@ -440,15 +424,8 @@ public interface TransformableWrappedStream<T, S extends TransformableWrappedStr
*/ */
default <R> EasyStream<R> flatMapIdx(final BiFunction<? super T, Integer, ? extends Stream<? extends R>> mapper) { default <R> EasyStream<R> flatMapIdx(final BiFunction<? super T, Integer, ? extends Stream<? extends R>> mapper) {
Objects.requireNonNull(mapper); Objects.requireNonNull(mapper);
if (isParallel()) { final MutableInt index = new MutableInt(NOT_FOUND_ELEMENT_INDEX);
final Map<Integer, T> idxMap = easyStream().toIdxMap(); return flatMap(e -> mapper.apply(e, index.incrementAndGet()));
return EasyStream.of(idxMap.entrySet())
.parallel(isParallel())
.flatMap(e -> mapper.apply(e.getValue(), e.getKey()));
} else {
final MutableInt index = new MutableInt(NOT_FOUND_ELEMENT_INDEX);
return flatMap(e -> mapper.apply(e, index.incrementAndGet()));
}
} }
/** /**
@ -553,15 +530,8 @@ public interface TransformableWrappedStream<T, S extends TransformableWrappedStr
*/ */
default <R> EasyStream<R> mapIdx(final BiFunction<? super T, Integer, ? extends R> mapper) { default <R> EasyStream<R> mapIdx(final BiFunction<? super T, Integer, ? extends R> mapper) {
Objects.requireNonNull(mapper); Objects.requireNonNull(mapper);
if (isParallel()) { final MutableInt index = new MutableInt(NOT_FOUND_ELEMENT_INDEX);
final Map<Integer, T> idxMap = easyStream().toIdxMap(); return map(e -> mapper.apply(e, index.incrementAndGet()));
return EasyStream.of(idxMap.entrySet())
.parallel(isParallel())
.map(e -> mapper.apply(e.getValue(), e.getKey()));
} else {
final MutableInt index = new MutableInt(NOT_FOUND_ELEMENT_INDEX);
return map(e -> mapper.apply(e, index.incrementAndGet()));
}
} }
/** /**

View File

@ -12,7 +12,6 @@ import org.junit.jupiter.api.Test;
import java.math.BigDecimal; import java.math.BigDecimal;
import java.math.RoundingMode; import java.math.RoundingMode;
import java.util.*; import java.util.*;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.function.Function; import java.util.function.Function;
import java.util.stream.Collectors; import java.util.stream.Collectors;
import java.util.stream.Stream; import java.util.stream.Stream;
@ -145,9 +144,6 @@ public class EasyStreamTest {
final List<String> list = Arrays.asList("dromara", "hutool", "sweet"); final List<String> list = Arrays.asList("dromara", "hutool", "sweet");
final List<String> mapIndex = EasyStream.of(list).mapIdx((e, i) -> i + 1 + "." + e).toList(); final List<String> mapIndex = EasyStream.of(list).mapIdx((e, i) -> i + 1 + "." + e).toList();
Assertions.assertEquals(Arrays.asList("1.dromara", "2.hutool", "3.sweet"), mapIndex); Assertions.assertEquals(Arrays.asList("1.dromara", "2.hutool", "3.sweet"), mapIndex);
// 并行流时正常
Assertions.assertEquals(Arrays.asList("1.dromara", "2.hutool", "3.sweet"),
EasyStream.of("dromara", "hutool", "sweet").parallel().mapIdx((e, i) -> i + 1 + "." + e).toList());
} }
@Test @Test
@ -199,10 +195,6 @@ public class EasyStreamTest {
final EasyStream.Builder<String> builder = EasyStream.builder(); final EasyStream.Builder<String> builder = EasyStream.builder();
EasyStream.of(list).forEachIdx((e, i) -> builder.accept(i + 1 + "." + e)); EasyStream.of(list).forEachIdx((e, i) -> builder.accept(i + 1 + "." + e));
Assertions.assertEquals(Arrays.asList("1.dromara", "2.hutool", "3.sweet"), builder.build().toList()); Assertions.assertEquals(Arrays.asList("1.dromara", "2.hutool", "3.sweet"), builder.build().toList());
// 并行流时正常
final AtomicInteger total = new AtomicInteger(0);
EasyStream.of("dromara", "hutool", "sweet").parallel().forEachIdx((e, i) -> total.addAndGet(i));
Assertions.assertEquals(3, total.get());
} }
@Test @Test
@ -223,10 +215,6 @@ public class EasyStreamTest {
final List<String> list = Arrays.asList("dromara", "hutool", "sweet"); final List<String> list = Arrays.asList("dromara", "hutool", "sweet");
final List<String> mapIndex = EasyStream.of(list).flatMapIdx((e, i) -> EasyStream.of(i + 1 + "." + e)).toList(); final List<String> mapIndex = EasyStream.of(list).flatMapIdx((e, i) -> EasyStream.of(i + 1 + "." + e)).toList();
Assertions.assertEquals(Arrays.asList("1.dromara", "2.hutool", "3.sweet"), mapIndex); Assertions.assertEquals(Arrays.asList("1.dromara", "2.hutool", "3.sweet"), mapIndex);
// 并行流时正常
Assertions.assertEquals(Arrays.asList("1.dromara", "2.hutool", "3.sweet"),
EasyStream.of("dromara", "hutool", "sweet").parallel()
.flatMapIdx((e, i) -> EasyStream.of(i + 1 + "." + e)).toList());
} }
@Test @Test
@ -280,9 +268,6 @@ public class EasyStreamTest {
final List<String> list = Arrays.asList("dromara", "hutool", "sweet"); final List<String> list = Arrays.asList("dromara", "hutool", "sweet");
final List<String> filterIndex = EasyStream.of(list).filterIdx((e, i) -> i < 2).toList(); final List<String> filterIndex = EasyStream.of(list).filterIdx((e, i) -> i < 2).toList();
Assertions.assertEquals(Arrays.asList("dromara", "hutool"), filterIndex); Assertions.assertEquals(Arrays.asList("dromara", "hutool"), filterIndex);
// 并行流时正常
Assertions.assertEquals(Arrays.asList("dromara", "hutool"),
EasyStream.of("dromara", "hutool", "sweet").parallel().filterIdx((e, i) -> i < 2).toList());
} }
@Test @Test
@ -350,7 +335,6 @@ public class EasyStreamTest {
public void testFindFirstIdx() { public void testFindFirstIdx() {
final List<Integer> list = Arrays.asList(null, 2, 3); final List<Integer> list = Arrays.asList(null, 2, 3);
Assertions.assertEquals(1, EasyStream.of(list).findFirstIdx(Objects::nonNull)); Assertions.assertEquals(1, EasyStream.of(list).findFirstIdx(Objects::nonNull));
Assertions.assertEquals(-1, (Object) EasyStream.of(list).parallel().findFirstIdx(Objects::nonNull));
} }
@Test @Test
@ -370,7 +354,6 @@ public class EasyStreamTest {
public void testFindLastIdx() { public void testFindLastIdx() {
final List<Integer> list = Arrays.asList(1, null, 3); final List<Integer> list = Arrays.asList(1, null, 3);
Assertions.assertEquals(2, (Object) EasyStream.of(list).findLastIdx(Objects::nonNull)); Assertions.assertEquals(2, (Object) EasyStream.of(list).findLastIdx(Objects::nonNull));
Assertions.assertEquals(-1, (Object) EasyStream.of(list).parallel().findLastIdx(Objects::nonNull));
} }
@Test @Test