package stream; import function.*; import java.util.HashSet; import java.util.Set; /** * @Author xiongyx * @Date 2019/3/6 * * stream实现 */ public class MyStream implements Stream { /** * 流的头部 * */ private T head; /** * 流的下一项求值函数 * */ private NextItemEvalProcess nextItemEvalProcess; /** * 是否是流的结尾 * */ private boolean isEnd; public static class Builder{ private MyStream target; public Builder() { this.target = new MyStream<>(); } public Builder head(T head){ target.head = head; return this; } Builder isEnd(boolean isEnd){ target.isEnd = isEnd; return this; } public Builder nextItemEvalProcess(NextItemEvalProcess nextItemEvalProcess){ target.nextItemEvalProcess = nextItemEvalProcess; return this; } public MyStream build(){ return target; } } //=================================API接口实现============================== @Override public MyStream map(Function mapper) { NextItemEvalProcess lastNextItemEvalProcess = this.nextItemEvalProcess; this.nextItemEvalProcess = new NextItemEvalProcess( ()->{ MyStream myStream = lastNextItemEvalProcess.eval(); return map(mapper, myStream); } ); // 求值链条 加入一个新的process map return new MyStream.Builder() .nextItemEvalProcess(this.nextItemEvalProcess) .build(); } @Override public MyStream flatMap(Function,T> mapper) { NextItemEvalProcess lastNextItemEvalProcess = this.nextItemEvalProcess; this.nextItemEvalProcess = new NextItemEvalProcess( ()->{ MyStream myStream = lastNextItemEvalProcess.eval(); return flatMap(mapper, Stream.makeEmptyStream(), myStream); } ); // 求值链条 加入一个新的process map return new MyStream.Builder() .nextItemEvalProcess(this.nextItemEvalProcess) .build(); } @Override public MyStream filter(Predicate predicate) { NextItemEvalProcess lastNextItemEvalProcess = this.nextItemEvalProcess; this.nextItemEvalProcess = new NextItemEvalProcess( ()-> { MyStream myStream = lastNextItemEvalProcess.eval(); return filter(predicate, myStream); } ); // 求值链条 加入一个新的process filter return this; } @Override public MyStream limit(int n) { NextItemEvalProcess lastNextItemEvalProcess = this.nextItemEvalProcess; this.nextItemEvalProcess = new NextItemEvalProcess( ()-> { MyStream myStream = lastNextItemEvalProcess.eval(); return limit(n, myStream); } ); // 求值链条 加入一个新的process limit return this; } @Override public MyStream distinct() { NextItemEvalProcess lastNextItemEvalProcess = this.nextItemEvalProcess; this.nextItemEvalProcess = new NextItemEvalProcess( ()-> { MyStream myStream = lastNextItemEvalProcess.eval(); return distinct(new HashSet<>(), myStream); } ); // 求值链条 加入一个新的process limit return this; } @Override public MyStream peek(ForEach consumer) { NextItemEvalProcess lastNextItemEvalProcess = this.nextItemEvalProcess; this.nextItemEvalProcess = new NextItemEvalProcess( ()-> { MyStream myStream = lastNextItemEvalProcess.eval(); return peek(consumer,myStream); } ); // 求值链条 加入一个新的process peek return this; } @Override public void forEach(ForEach consumer) { // 终结操作 直接开始求值 forEach(consumer,this.eval()); } @Override public R reduce(R initVal, BiFunction accumulator) { // 终结操作 直接开始求值 return reduce(initVal,accumulator,this.eval()); } @Override public R collect(Collector collector) { // 终结操作 直接开始求值 A result = collect(collector,this.eval()); // 通过finish方法进行收尾 return collector.finisher().apply(result); } @Override public T max(Comparator comparator) { // 终结操作 直接开始求值 MyStream eval = this.eval(); if(eval.isEmptyStream()){ return null; }else{ return max(comparator,eval,eval.head); } } @Override public T min(Comparator comparator) { // 终结操作 直接开始求值 MyStream eval = this.eval(); if(eval.isEmptyStream()){ return null; }else{ return min(comparator,eval,eval.head); } } @Override public int count() { // 终结操作 直接开始求值 return count(this.eval(),0); } @Override public boolean anyMatch(Predicate predicate) { // 终结操作 直接开始求值 return anyMatch(predicate,this.eval()); } @Override public boolean allMatch(Predicate predicate) { // 终结操作 直接开始求值 return allMatch(predicate,this.eval()); } //===============================私有方法==================================== /** * 递归函数 配合API.map * */ private static MyStream map(Function mapper, MyStream myStream){ if(myStream.isEmptyStream()){ return Stream.makeEmptyStream(); } R head = mapper.apply(myStream.head); return new MyStream.Builder() .head(head) .nextItemEvalProcess(new NextItemEvalProcess(()->map(mapper, myStream.eval()))) .build(); } /** * 递归函数 配合API.flatMap * */ private static MyStream flatMap(Function,T> mapper, MyStream headMyStream, MyStream myStream){ if(headMyStream.isEmptyStream()){ if(myStream.isEmptyStream()){ return Stream.makeEmptyStream(); }else{ T outerHead = myStream.head; MyStream newHeadMyStream = mapper.apply(outerHead); return flatMap(mapper, newHeadMyStream.eval(), myStream.eval()); } }else{ return new MyStream.Builder() .head(headMyStream.head) .nextItemEvalProcess(new NextItemEvalProcess(()-> flatMap(mapper, headMyStream.eval(), myStream))) .build(); } } /** * 递归函数 配合API.filter * */ private static MyStream filter(Predicate predicate, MyStream myStream){ if(myStream.isEmptyStream()){ return Stream.makeEmptyStream(); } if(predicate.satisfy(myStream.head)){ return new Builder() .head(myStream.head) .nextItemEvalProcess(new NextItemEvalProcess(()->filter(predicate, myStream.eval()))) .build(); }else{ return filter(predicate, myStream.eval()); } } /** * 递归函数 配合API.limit * */ private static MyStream limit(int num, MyStream myStream){ if(num == 0 || myStream.isEmptyStream()){ return Stream.makeEmptyStream(); } return new MyStream.Builder() .head(myStream.head) .nextItemEvalProcess(new NextItemEvalProcess(()->limit(num-1, myStream.eval()))) .build(); } /** * 递归函数 配合API.distinct * */ private static MyStream distinct(Set distinctSet,MyStream myStream){ if(myStream.isEmptyStream()){ return Stream.makeEmptyStream(); } if(!distinctSet.contains(myStream.head)){ // 加入集合 distinctSet.add(myStream.head); return new Builder() .head(myStream.head) .nextItemEvalProcess(new NextItemEvalProcess(()->distinct(distinctSet, myStream.eval()))) .build(); }else{ return distinct(distinctSet, myStream.eval()); } } /** * 递归函数 配合API.peek * */ private static MyStream peek(ForEach consumer,MyStream myStream){ if(myStream.isEmptyStream()){ return Stream.makeEmptyStream(); } consumer.apply(myStream.head); return new MyStream.Builder() .head(myStream.head) .nextItemEvalProcess(new NextItemEvalProcess(()->peek(consumer, myStream.eval()))) .build(); } /** * 递归函数 配合API.forEach * */ private static void forEach(ForEach consumer, MyStream myStream){ if(myStream.isEmptyStream()){ return; } consumer.apply(myStream.head); forEach(consumer, myStream.eval()); } /** * 递归函数 配合API.reduce * */ private static R reduce(R initVal, BiFunction accumulator, MyStream myStream){ if(myStream.isEmptyStream()){ return initVal; } T head = myStream.head; R result = reduce(initVal,accumulator, myStream.eval()); return accumulator.apply(result,head); } /** * 递归函数 配合API.collect * */ private static A collect(Collector collector, MyStream myStream){ if(myStream.isEmptyStream()){ return collector.supplier().get(); } T head = myStream.head; A tail = collect(collector, myStream.eval()); return collector.accumulator().apply(tail,head); } /** * 递归函数 配合API.max * */ private static T max(Comparator comparator, MyStream myStream, T max){ if(myStream.isEnd){ return max; } T head = myStream.head; // head 和 max 进行比较 if(comparator.compare(head,max) > 0){ // head 较大 作为新的max传入 return max(comparator, myStream.eval(),head); }else{ // max 较大 不变 return max(comparator, myStream.eval(),max); } } /** * 递归函数 配合API.min * */ private static T min(Comparator comparator, MyStream myStream, T min){ if(myStream.isEnd){ return min; } T head = myStream.head; // head 和 min 进行比较 if(comparator.compare(head,min) < 0){ // head 较小 作为新的min传入 return min(comparator, myStream.eval(),head); }else{ // min 较小 不变 return min(comparator, myStream.eval(),min); } } /** * 递归函数 配合API.count * */ private static int count(MyStream myStream, int count){ if(myStream.isEmptyStream()){ return count; } // count+1 进行递归 return count(myStream.eval(),count+1); } /** * 递归函数 配合API.anyMatch * */ private static boolean anyMatch(Predicate predicate,MyStream myStream){ if(myStream.isEmptyStream()){ // 截止末尾,不存在任何匹配项 return false; } // 谓词判断 if(predicate.satisfy(myStream.head)){ // 匹配 存在匹配项 返回true return true; }else{ // 不匹配,继续检查,直到存在匹配项 return anyMatch(predicate,myStream.eval()); } } /** * 递归函数 配合API.anyMatch * */ private static boolean allMatch(Predicate predicate,MyStream myStream){ if(myStream.isEmptyStream()){ // 全部匹配 return true; } // 谓词判断 if(predicate.satisfy(myStream.head)){ // 当前项匹配,继续检查 return allMatch(predicate,myStream.eval()); }else{ // 存在不匹配的项,返回false return false; } } /** * 当前流强制求值 * @return 求值之后返回一个新的流 * */ private MyStream eval(){ return this.nextItemEvalProcess.eval(); } /** * 当前流 为空 * */ private boolean isEmptyStream(){ return this.isEnd; } @Override public String toString() { return "MyStream{" + "head=" + head + ", isEnd=" + isEnd + '}'; } }