玩转 Java 8 Stream API

image-20211031204304133

先贴上几个案例,水平高超的同学可以挑战一下:

  1. 从员工集合中筛选出salary大于8000的员工,并放置到新的集合里。
  2. 统计员工的最高薪资、平均薪资、薪资之和。
  3. 将员工按薪资从高到低排序,同样薪资者年龄小者在前。
  4. 将员工按性别分类,将员工按性别和地区分类,将员工按薪资是否高于8000分为两部分。

用传统的迭代处理也不是很难,但代码就显得冗余了,跟Stream相比高下立判。

1 Stream概述

Java 8 是一个非常成功的版本,这个版本新增的Stream,配合同版本出现的 Lambda ,给我们操作集合(Collection)提供了极大的便利。

那么什么是Stream

Stream将要处理的元素集合看作一种流,在流的过程中,借助Stream API对流中的元素进行操作,比如:筛选、排序、聚合等。

Stream可以由数组或集合创建,对流的操作分为两种:

  1. 中间操作,每次返回一个新的流,可以有多个。
  2. 终端操作,每个流只能进行一次终端操作,终端操作结束后流无法再次使用。终端操作会产生一个新的集合或值。

另外,Stream有几个特性:

  1. stream不存储数据,而是按照特定的规则对数据进行计算,一般会输出结果。
  2. stream不会改变数据源,通常情况下会产生一个新的集合或一个值。
  3. stream具有延迟执行特性,只有调用终端操作时,中间操作才会执行。

2 Stream的创建

Stream可以通过集合数组创建。

1、通过 java.util.Collection.stream() 方法用集合创建流

List<String> list = Arrays.asList("a", "b", "c");
// 创建一个顺序流
Stream<String> stream = list.stream();
// 创建一个并行流
Stream<String> parallelStream = list.parallelStream();

2、使用java.util.Arrays.stream(T[] array)方法用数组创建流

int[] array={1,3,5,6,8};
IntStream stream = Arrays.stream(array);

3、使用Stream的静态方法:of()、iterate()、generate()

Stream<Integer> stream = Stream.of(1, 2, 3, 4, 5, 6);

Stream<Integer> stream2 = Stream.iterate(0, (x) -> x + 3).limit(4);
stream2.forEach(System.out::println); // 0 2 4 6 8 10

Stream<Double> stream3 = Stream.generate(Math::random).limit(3);
stream3.forEach(System.out::println);

输出结果:

0 3 6 9

0.6796156909271994 0.1914314208854283 0.8116932592396652

streamparallelStream的简单区分: stream是顺序流,由主线程按顺序对流执行操作,而parallelStream是并行流,内部以多线程并行执行的方式对流进行操作,但前提是流中的数据处理没有顺序要求。例如筛选集合中的奇数,两者的处理不同之处:

image-20211031203335805

如果流中的数据量足够大,并行流可以加快处速度。

除了直接创建并行流,还可以通过parallel()把顺序流转换成并行流:

Optional<Integer> findFirst = list.stream().parallel().filter(x->x>6).findFirst();

3 Stream的使用

在使用stream之前,先理解一个概念:Optional

Optional类是一个可以为null的容器对象。如果值存在则isPresent()方法会返回true,调用get()方法会返回该对象, 否则会抛出异常。Optional类还可以用更优雅的方式进行判空处理, 更详细说明请见:https://docs.oracle.com/javase/8/docs/api/java/util/Optional.html

接下来,大批代码向你袭来!我将用20个案例将Stream的使用整得明明白白,只要跟着敲一遍代码,就能很好地掌握。

案例使用的员工类

这是案例中使用的员工类:

public class Person {
    private String name;  // 姓名
    private int salary; // 薪资
    private int age; // 年龄
    private String sex; //性别
    private String area;  // 地区

    // 构造方法
    public Person(String name, int salary, String sex, String area) {
        this.name = name;
        this.salary = salary;
        this.age = age;
        this.sex = sex;
        this.area = area;
    }

    public Person(String name, int salary, int age, String sex, String area) {
        this.name = name;
        this.salary = salary;
        this.age = age;
        this.sex = sex;
        this.area = area;
    }

    // 省略了get和set,请自行添加
    public String getName() {
        return name;
    }

    public void setName(String name) {
        this.name = name;
    }

    public int getSalary() {
        return salary;
    }

    public void setSalary(int salary) {
        this.salary = salary;
    }

    public int getAge() {
        return age;
    }

    public void setAge(int age) {
        this.age = age;
    }

    public String getSex() {
        return sex;
    }

    public void setSex(String sex) {
        this.sex = sex;
    }

    public String getArea() {
        return area;
    }

    public void setArea(String area) {
        this.area = area;
    }
}

3.1 遍历/匹配(foreach/find/match)

Stream也是支持类似集合的遍历和匹配元素的,只是Stream中的元素是以Optional类型存在的。Stream的遍历、匹配非常简单。

image-20211031203355248

// 遍历输出符合条件的元素
list.stream().filter(x -> x > 6).forEach(System.out::println);
// 匹配第一个
Optional<Integer> findFirst = list.stream().filter(x -> x > 6).findFirst();
// 匹配任意(适用于并行流)
Optional<Integer> findAny = list.parallelStream().filter(x -> x > 6).findAny();
// 是否包含符合特定条件的元素
System.out.println("匹配第一个值:" + findFirst.get());
System.out.println("匹配任意一个值:" + findAny.get());
boolean anyMatch = list.stream().anyMatch(x -> x > 6);
System.out.println("是否存在大于6的值:" + anyMatch);

3.2 筛选(filter)

筛选,是按照一定的规则校验流中的元素,将符合条件的元素提取到新的流中的操作。

image-20211031210320051

案例一:筛选出Integer集合中大于7的元素,并打印出来

List<Integer> list2 = Arrays.asList(6, 7, 3, 8, 1, 2, 9);
Stream<Integer> stream = list2.stream();
stream.filter(x -> x > 7).forEach(System.out::println);

预期结果:

8 9

案例二:筛选员工中工资高于8000的人,并形成新的集合。 形成新集合依赖collect(收集),后文有详细介绍。

List<Person> personList = new ArrayList<Person>();
personList.add(new Person("Tom", 8900, "male", "New York"));
personList.add(new Person("Jack", 7000, "male", "Washington"));
personList.add(new Person("Lily", 7800, "female", "Washington"));
personList.add(new Person("Anni", 8200, "female", "New York"));
personList.add(new Person("Owen", 9500, "male", "New York"));
personList.add(new Person("Alisa", 7900, "female", "New York"));

// 筛选员工中工资高于8000的人,并形成新的集合
List<String> fiterList = personList.stream().filter(x -> x.getSalary() > 8000).map(Person::getName).collect(Collectors.toList());
System.out.print("高于8000的员工姓名:" + fiterList);

运行结果:

高于8000的员工姓名:[Tom, Anni, Owen]

3.3 聚合(max/min/count)

maxmincount这些字眼你一定不陌生,没错,在mysql中我们常用它们进行数据统计。Java stream 中也引入了这些概念和用法,极大地方便了我们对集合、数组的数据统计工作。

image-20211031203425525

案例一:获取String集合中最长的元素。

// 获取String集合中最长的元素
List<String> list = Arrays.asList("adnm", "admmt", "pot", "xbangd", "weoujgsd");
Optional<String> max = list.stream().max(Comparator.comparing(String::length));
System.out.println("最长的字符串:" + max.get());

输出结果:

最长的字符串:weoujgsd

案例二:获取Integer集合中的最大值。

List<Integer> list = Arrays.asList(7, 6, 9, 4, 11, 6);

// 自然排序
Optional<Integer> max = list.stream().max(Integer::compareTo);
Optional<Integer> max1 = list.stream().max(Comparator.naturalOrder());
System.out.println("自然排序的最大值:" + max.get());
System.out.println("自然排序的最大值:" + max1.get());

// 自定义排序
Optional<Integer> max2 = list.stream().max((o1, o2) -> o1.compareTo(o2));
System.out.println("自定义排序的最大值:" + max2.get());

输出结果:

自然排序的最大值:11

自定义排序的最大值:11

案例三:获取员工工资最高的人。

List<Person> personList = new ArrayList<>();
personList.add(new Person("Tom", 8900, "male", "New York"));
personList.add(new Person("Jack", 7000, "male", "Washington"));
personList.add(new Person("Lily", 7800, "female", "Washington"));
personList.add(new Person("Anni", 8200, "female", "New York"));
personList.add(new Person("Owen", 9500, "male", "New York"));
personList.add(new Person("Alisa", 7900, "female", "New York"));

Optional<Person> max = personList.stream().max(Comparator.comparingInt(Person::getSalary));
System.out.println("员工工资最大值:" + max.get().getSalary());

输出结果:

员工工资最大值:9500

案例四:计算Integer集合中大于6的元素的个数。

List<Integer> list = Arrays.asList(7, 6, 4, 8, 2, 11, 9);
long count = list.stream().filter(x -> x > 6).count();
System.out.println("list中大于6的元素个数:" + count);

输出结果:

list中大于6的元素个数:4

3.4 映射(map/flatMap)

映射,可以将一个流的元素按照一定的映射规则映射到另一个流中。分为mapflatMap

  • map:接收一个函数作为参数,该函数会被应用到每个元素上,并将其映射成一个新的元素。
  • flatMap:接收一个函数作为参数,将流中的每个值都换成另一个流,然后把所有流连接成一个流。

image-20211031203517018

案例一:英文字符串数组的元素全部改为大写。整数数组每个元素+3。

String[] strArr = { "abcd", "bcdd", "defde", "fTr" };
List<String> strList = Arrays.stream(strArr).map(String::toUpperCase).collect(Collectors.toList());
System.out.println("每个元素大写:" + strList);

List<Integer> intList = Arrays.asList(1, 3, 5, 7, 9, 11);
List<Integer> intListNew = intList.stream().map(x -> x + 3).collect(Collectors.toList());
System.out.println("每个元素+3:" + intListNew);

输出结果:

每个元素大写:[ABCD, BCDD, DEFDE, FTR]

每个元素+3:[4, 6, 8, 10, 12, 14]

案例二:将员工的薪资全部增加1000。

List<Person> personList = new ArrayList<Person>();
personList.add(new Person("Tom", 8900, "male", "New York"));
personList.add(new Person("Jack", 7000, "male", "Washington"));
personList.add(new Person("Lily", 7800, "female", "Washington"));
personList.add(new Person("Anni", 8200, "female", "New York"));
personList.add(new Person("Owen", 9500, "male", "New York"));
personList.add(new Person("Alisa", 7900, "female", "New York"));

// 不改变原来员工集合的方式
List<Person> personListNew = personList.stream().map(person -> {
    Person personNew = new Person(person.getName(), 0, null, null);
    personNew.setSalary(person.getSalary() + 10000);
    return personNew;
}).collect(Collectors.toList());
System.out.println("一次改动前:" + personList.get(0).getName() + "-->" + personList.get(0).getSalary());
System.out.println("一次改动后:" + personListNew.get(0).getName() + "-->" + personListNew.get(0).getSalary());

// 改变原来员工集合的方式
List<Person> personListNew2 = personList.stream().map(person -> {
    person.setSalary(person.getSalary() + 10000);
    return person;
}).collect(Collectors.toList());
System.out.println("二次改动前:" + personList.get(0).getName() + "-->" + personListNew.get(0).getSalary());
System.out.println("二次改动后:" + personListNew2.get(0).getName() + "-->" + personListNew.get(0).getSalary());

输出结果:

一次改动前:Tom–>8900

一次改动后:Tom–>18900

二次改动前:Tom–>18900

二次改动后:Tom–>18900

案例三:将两个字符数组合并成一个新的字符数组。

List<String> list = Arrays.asList("m,k,l,a", "1,3,5,7");
List<String> listNew = list.stream().flatMap(s -> {
    // 将每个元素转换成一个stream
    String[] split = s.split(",");
    Stream<String> s2 = Arrays.stream(split);
    return s2;
}).collect(Collectors.toList());

System.out.println("处理前的集合:" + list);
System.out.println("处理后的集合:" + listNew);

输出结果:

处理前的集合:[m-k-l-a, 1-3-5]

处理后的集合:[m, k, l, a, 1, 3, 5, 7]

3.5 归约(reduce)

归约,也称缩减, 其实就是从前往后两两归并, 最后得到一个总的归并的结果,从结果来看是把一个流缩减成一个值,能实现对集合求和、求乘积和求最值操作。

image-20211031203558850

案例一:求Integer集合的元素之和、乘积和最大值。

List<Integer> list = Arrays.asList(1, 3, 2, 8, 11, 4);
// 求和方式1
Optional<Integer> sum = list.stream().reduce((x, y) -> x + y);
// 求和方式2
Optional<Integer> sum2 = list.stream().reduce(Integer::sum);
// 求和方式3 - 第一个参数是第一次用于累加的数
Integer sum3 = list.stream().reduce(0, Integer::sum);
System.out.println("list求和:" + sum.get() + "," + sum2.get() + "," + sum3);

// 求乘积
Optional<Integer> product = list.stream().reduce((x, y) -> x * y);
System.out.println("list求积:" + product.get());

// 求最大值方式1
Optional<Integer> max = list.stream().reduce((x, y) -> x > y ? x : y);
// 求最大值写法2 - 第一个参数是第一次用于比较的数
Integer max2 = list.stream().reduce(Integer.MIN_VALUE, Integer::max);
System.out.println("list求最大值:" + max.get() + "," + max2);

输出结果:

list求和:29,29,29

list求积:2112 list

list求最大值:11,11

案例二:求所有员工的工资之和和最高工资。

List<Person> personList = new ArrayList<Person>();
personList.add(new Person("Tom", 8900, "male", "New York"));
personList.add(new Person("Jack", 7000, "male", "Washington"));
personList.add(new Person("Lily", 7800, "female", "Washington"));
personList.add(new Person("Anni", 8200, "female", "New York"));
personList.add(new Person("Owen", 9500, "male", "New York"));
personList.add(new Person("Alisa", 7900, "female", "New York"));

// 求工资之和方式1:
Optional<Integer> sumSalary = personList.stream().map(Person::getSalary).reduce(Integer::sum);
// 求工资之和方式2:
Integer sumSalary2 = personList.stream().reduce(0, (first, second) -> first += second.getSalary(),
        (sum1, sum2) -> sum1 + sum2);
// 求工资之和方式3:
Integer sumSalary3 = personList.stream().reduce(0, (first, second) -> first += second.getSalary(), Integer::sum);
System.out.println("工资之和:" + sumSalary.get() + "," + sumSalary2 + "," + sumSalary3);

// 求最高工资方式1:
Integer maxSalary = personList.stream().reduce(0, (max, p) -> max > p.getSalary() ? max : p.getSalary(),
        Integer::max);
// 求最高工资方式2:
Integer maxSalary2 = personList.stream().reduce(0, (max, p) -> max > p.getSalary() ? max : p.getSalary(),
        (max1, max2) -> max1 > max2 ? max1 : max2);
System.out.println("最高工资:" + maxSalary + "," + maxSalary2);

输出结果:

工资之和:49300,49300,49300

最高工资:9500,9500

3.6 收集(collect)

collect,收集,可以说是内容最繁多、功能最丰富的部分了。从字面上去理解,就是把一个流收集起来,最终可以是收集成一个值也可以收集成一个新的集合。

collect主要依赖java.util.stream.Collectors类内置的静态方法。

3.6.1 归集(toList/toSet/toMap)

因为流不存储数据,那么在流中的数据完成处理后,需要将流中的数据重新归集到新的集合里。toListtoSettoMap比较常用,另外还有toCollectiontoConcurrentMap等复杂一些的用法。

下面用一个案例演示toListtoSettoMap

List<Integer> list = Arrays.asList(1, 6, 3, 4, 6, 7, 9, 6, 20);
List<Integer> listNew = list.stream().filter(x -> x % 2 == 0).collect(Collectors.toList());
System.out.println("toList:" + listNew);

Set<Integer> set = list.stream().filter(x -> x % 2 == 0).collect(Collectors.toSet());
System.out.println("toSet:" + set);

List<Person> personList = new ArrayList<Person>();
personList.add(new Person("Tom", 8900, "male", "New York"));
personList.add(new Person("Jack", 7000, "male", "Washington"));
personList.add(new Person("Lily", 7800, "female", "Washington"));
personList.add(new Person("Anni", 8200, "female", "New York"));
Map<?, Person> map = personList.stream().filter(p -> p.getSalary() > 8000)
        .collect(Collectors.toMap(Person::getName, p -> p));
System.out.println("toMap:" + map);

运行结果:

toList:[6, 4, 6, 6, 20]

toSet:[4, 20, 6]

toMap:{Tom=mutest.Person@5fd0d5ae, Anni=mutest.Person@2d98a335}

3.6.2 统计(count/averaging)

Collectors提供了一系列用于数据统计的静态方法:

  • 计数:count
  • 平均值:averagingIntaveragingLongaveragingDouble
  • 最值:maxByminBy
  • 求和:summingIntsummingLongsummingDouble
  • 统计以上所有:summarizingIntsummarizingLongsummarizingDouble

案例:统计员工人数、平均工资、工资总额、最高工资。

List<Person> personList = new ArrayList<Person>();
personList.add(new Person("Tom", 8900, "male", "New York"));
personList.add(new Person("Jack", 7000, "male", "Washington"));
personList.add(new Person("Lily", 7800, "female", "Washington"));

// 求总数
// Long count = (long) personList.size();
// Long count = personList.stream().count();
Long count = personList.stream().collect(Collectors.counting());
System.out.println("员工总数:" + count);

// 求平均工资
Double average = personList.stream().collect(Collectors.averagingDouble(Person::getSalary));
System.out.println("员工平均工资:" + average);

// 求最高工资
Optional<Integer> max = personList.stream().map(Person::getSalary).collect(Collectors.maxBy(Integer::compare));
System.out.println("员工最高工资:" + max);

// 求工资之和
Integer sum = personList.stream().collect(Collectors.summingInt(Person::getSalary));
System.out.println("员工工资总和:" + sum);

// 一次性统计所有信息
DoubleSummaryStatistics collect = personList.stream().collect(Collectors.summarizingDouble(Person::getSalary));
System.out.println("员工工资所有统计:" + collect);

运行结果:

员工总数:3 员工平均工资:7900.0 员工工资总和:23700 员工工资所有统计:DoubleSummaryStatistics{count=3, sum=23700.000000,min=7000.000000, average=7900.000000, max=8900.000000}

3.6.3 分组(partitioningBy/groupingBy)

  • 分区:将stream按条件分为两个Map,比如员工按薪资是否高于8000分为两部分。
  • 分组:将集合分为多个Map,比如员工按性别分组。有单级分组和多级分组。

image-20211031203705460

案例:将员工按薪资是否高于8000分为两部分;将员工按性别和地区分组

List<Person> personList = new ArrayList<Person>();
personList.add(new Person("Tom", 8900, "male", "New York"));
personList.add(new Person("Jack", 7000, "male", "Washington"));
personList.add(new Person("Lily", 7800, "female", "Washington"));
personList.add(new Person("Anni", 8200, "female", "New York"));
personList.add(new Person("Owen", 9500, "male", "New York"));
personList.add(new Person("Alisa", 7900, "female", "New York"));

// 将员工按薪资是否高于8000分组
Map<Boolean, List<Person>> part = personList.stream().collect(Collectors.partitioningBy(x -> x.getSalary() > 8000));
// 将员工按性别分组
Map<String, List<Person>> group = personList.stream().collect(Collectors.groupingBy(Person::getSex));
// 将员工先按性别分组,再按地区分组
Map<String, Map<String, List<Person>>> group2 = personList.stream().collect(Collectors.groupingBy(Person::getSex, Collectors.groupingBy(Person::getArea)));
System.out.println("员工按薪资是否大于8000分组情况:" + part);
System.out.println("员工按性别分组情况:" + group);
System.out.println("员工按性别、地区:" + group2);

输出结果:

员工按薪资是否大于8000分组情况:{false=[mutest.Person@2d98a335, mutest.Person@16b98e56, mutest.Person@7ef20235], true=[mutest.Person@27d6c5e0, mutest.Person@4f3f5b24, mutest.Person@15aeb7ab]}  

员工按性别分组情况:{female=[mutest.Person@16b98e56, mutest.Person@4f3f5b24, mutest.Person@7ef20235], male=[mutest.Person@27d6c5e0, mutest.Person@2d98a335, mutest.Person@15aeb7ab]}  

员工按性别、地区:{female={New York=[mutest.Person@4f3f5b24, mutest.Person@7ef20235], Washington=[mutest.Person@16b98e56]}, male={New York=[mutest.Person@27d6c5e0, mutest.Person@15aeb7ab], Washington=[mutest.Person@2d98a335]}}  

3.6.4 接合(joining)

joining可以将stream中的元素用特定的连接符(没有的话,则直接连接)连接成一个字符串。

List<Person> personList = new ArrayList<Person>();
personList.add(new Person("Tom", 8900, 23, "male", "New York"));
personList.add(new Person("Jack", 7000, 25, "male", "Washington"));
personList.add(new Person("Lily", 7800, 21, "female", "Washington"));

String names = personList.stream().map(p -> p.getName()).collect(Collectors.joining(","));
System.out.println("所有员工的姓名:" + names);
List<String> list = Arrays.asList("A", "B", "C");
String string = list.stream().collect(Collectors.joining("-"));
System.out.println("拼接后的字符串:" + string);

运行结果:

所有员工的姓名:Tom,Jack,Lily 拼接后的字符串:A-B-C

3.6.5 集合工具类的归约方法(reducing)

stream本身的reduce方法也可以替换成Collectors类提供的reducing`方法。

List<Person> personList = new ArrayList<>();
personList.add(new Person("Tom", 8900, 23, "male", "New York"));
personList.add(new Person("Jack", 7000, 25, "male", "Washington"));
personList.add(new Person("Lily", 7800, 21, "female", "Washington"));

// 每个员工减去起征点后的薪资之和(这个例子并不严谨,但一时没想到好的例子)
Integer sum1 = personList.stream().collect(Collectors.reducing(0, Person::getSalary, (i, j) -> (i + j - 5000)));
Integer sum2 = personList.stream().collect(Collectors.reducing(0, Person::getSalary, (i, j) -> (i + j - 5000)));
System.out.println("员工扣税薪资总和:" + sum1);
System.out.println("员工扣税薪资总和:" + sum2);

// stream的reduce(建议)
Optional<Integer> sum3 = personList.stream().map(Person::getSalary).reduce(Integer::sum);
System.out.println("员工薪资总和:" + sum3.get());

运行结果:

员工扣税薪资总和:8700 员工薪资总和:23700

3.7 排序(sorted)

sorted,中间操作。有两种排序:

  • sorted():自然排序,流中元素需实现Comparable接口
  • sorted(Comparator com):Comparator排序器自定义排序

案例:将员工按工资由高到低(工资一样则按年龄由大到小)排序

List<Person> personList = new ArrayList<Person>();

personList.add(new Person("Sherry", 9000, 24, "female", "New York"));
personList.add(new Person("Tom", 8900, 22, "male", "Washington"));
personList.add(new Person("Jack", 9000, 25, "male", "Washington"));
personList.add(new Person("Lily", 8800, 26, "male", "New York"));
personList.add(new Person("Alisa", 9000, 26, "female", "New York"));

// 按工资升序排序(自然排序)
List<String> newList = personList.stream().sorted(Comparator.comparing(Person::getSalary)).map(Person::getName)
        .collect(Collectors.toList());
// 按工资倒序排序
List<String> newList2 = personList.stream().sorted(Comparator.comparing(Person::getSalary).reversed())
        .map(Person::getName).collect(Collectors.toList());
// 先按工资再按年龄升序排序
List<String> newList3 = personList.stream()
        .sorted(Comparator.comparing(Person::getSalary).thenComparing(Person::getAge)).map(Person::getName)
        .collect(Collectors.toList());
// 先按工资再按年龄自定义排序(降序)
List<String> newList4 = personList.stream().sorted((p1, p2) -> {
    if (p1.getSalary() == p2.getSalary()) {
        return p2.getAge() - p1.getAge(); // 降序
    } else {
        return p2.getSalary() - p1.getSalary(); // 降序
    }
}).map(Person::getName).collect(Collectors.toList());

System.out.println("按工资升序排序:" + newList);
System.out.println("按工资降序排序:" + newList2);
System.out.println("先按工资再按年龄升序排序:" + newList3);
System.out.println("先按工资再按年龄自定义降序排序:" + newList4);

运行结果:

按工资自然排序:[Lily, Tom, Sherry, Jack, Alisa] 按工资降序排序:[Sherry, Jack, Alisa,Tom, Lily] 先按工资再按年龄自然排序:[Sherry, Jack, Alisa, Tom, Lily] 先按工资再按年龄自定义降序排序:[Alisa, Jack, Sherry, Tom, Lily]

3.8 提取/组合

流也可以进行合并、去重、限制、跳过等操作。

image-20211031204447910

String[] arr1 = { "a", "b", "c", "d" };
String[] arr2 = { "d", "e", "f", "g" };

Stream<String> stream1 = Stream.of(arr1);
Stream<String> stream2 = Stream.of(arr2);

// concat:合并两个流 distinct:去重
List<String> newList = Stream.concat(stream1, stream2).distinct().collect(Collectors.toList());
System.out.println("concat and distinct:" + newList);

// limit:限制从流中获得前n个数据
// iterate: Returns an stream by a function to an initial element(seed)
List<Integer> collect = Stream.iterate(1, x -> x + 2).limit(10).collect(Collectors.toList());
System.out.println("limit:" + collect);

// skip:跳过前n个数据
List<Integer> collect2 = Stream.iterate(1, x -> x + 2).skip(1).limit(5).collect(Collectors.toList());
System.out.println("skip:" + collect2);

运行结果:

流合并:[a, b, c, d, e, f, g] limit:[1, 3, 5, 7, 9, 11, 13, 15, 17, 19] skip:[3, 5, 7, 9, 11]

原文链接及版权说明:

原文链接:https://blog.csdn.net/mu_wind/article/details/109516995

版权声明:本文为CSDN博主「云深i不知处」的原创文章,遵循CC 4.0 BY-SA版权协议,转载请附上原文出处链接及本声明。

并行流

对于CPU密集型任务使用并行流可以利用多线程提高执行效率.这里使用线程的sleep方法来模拟耗时的CPU密集型任务(经实验只有出现阻塞线程的操作才会使用多线程, 否则只会在主线程中进行计算), 这样实际计算就会启动ForkJoinPool中的worker线程来并发执行提高计算效率.

默认worker数量是CPU核数-1, 但是可以使用虚拟机选项-Djava.util.concurrent.ForkJoinPool.common.parallelism=N设置worker的数量为N(最大值为32767)。

    @Test
    public void loopingTest(){
        // 顺序流
        AtomicInteger result = new AtomicInteger();
        long start = System.nanoTime();
        for (int x = 1; x < 1000; x++) {
            Utils.sleep(10);
            result.addAndGet(x);
        }
        long end = System.nanoTime();
        System.out.printf("time spent for normal  looping: %.3f sec.%n",(end-start)*1E-9);
    }

    @Test
    public void sequenceStreamTest(){
        // 顺序流
        AtomicInteger result = new AtomicInteger();
        long start = System.nanoTime();
        IntStream.range(1, 1000).forEach(x->{
            Utils.sleep(10);
            result.addAndGet(x);
        });
        long end = System.nanoTime();
        System.out.printf("time spent for sequence stream: %.3f sec.%n",(end-start)*1E-9);
    }

    @Test
    public void parallelStreamTest(){
        // 并行流
        // 对于CPU密集型任务使用并行流可以提高执行效率
        // 此系统属性用来指定并行流计算所使用的 ForkJoinPool.commonPool-worker 的线程数量(并行度)
        System.setProperty("java.util.concurrent.ForkJoinPool.common.parallelism", "50");
        AtomicInteger result = new AtomicInteger();
        long start = System.nanoTime();
        IntStream.range(1, 1000).parallel().forEach(x->{
            Utils.sleep(10);
            result.addAndGet(x);
        });
        long end = System.nanoTime();
        System.out.printf("time spent for parallel stream: %.3f sec.%n",(end-start)*1E-9);
    }

输出结果:

time spent for parallel stream:  0.376 sec.
time spent for sequence stream: 15.917 sec.
time spent for normal  looping: 17.574 sec.

可以看到,当ForkJoinPool.commonPool-worker数量为50时, 原来顺使用序流需要运行17秒, 而并行流只需要1秒不到.

使用并行流实现WordCount

再看另一个例子, 分别使用顺序流和并行流(并行度设置为100)计算词频统计.

    @Test
    public void wordCountBySequenceStream() {

        List<String> lines;
        lines = Arrays.asList(
                "the cow jumped over the moon",
                "an apple a day keep the doctor away",
                "snow white and the seven dwarfs",
                "i am at two with nature",
                "the cow jumped over the moon",
                "an apple a day keep the doctor away",
                "snow white and the seven dwarfs",
                "i am at two with nature"
        );

        HashMap<String, Long> resultMap = new HashMap<>();

        long start = System.nanoTime();
        lines.stream().flatMap(line -> {
                    Utils.sleep(50);
                    return Arrays.stream(line.toLowerCase().split(" "));
                }
        ).forEach(word -> {
            Utils.sleep(50);
            if (resultMap.containsKey(word)) {
                resultMap.put(word, resultMap.get(word) + 1L);
            } else {
                resultMap.put(word, 1L);
            }
        });
        long end = System.nanoTime();
        System.out.printf("time spent for wordcount by sequence stream: %.3f sec.%n", (end - start) * 1E-9);

        resultMap.forEach(MxWordCountByParalleStreamTest::accept);
    }

    @Test
    public void wordCountByParallelStream() {
        System.setProperty("java.util.concurrent.ForkJoinPool.common.parallelism", "100");

        List<String> lines;
        lines = Arrays.asList(
                "the cow jumped over the moon",
                "an apple a day keep the doctor away",
                "snow white and the seven dwarfs",
                "i am at two with nature",
                "the cow jumped over the moon",
                "an apple a day keep the doctor away",
                "snow white and the seven dwarfs",
                "i am at two with nature"
        );

        HashMap<String, Long> resultMap = new HashMap<>();

        long start = System.nanoTime();
        lines.parallelStream().flatMap(line -> {
            Utils.sleep(50);
            return Arrays.stream(line.toLowerCase().split(" "));
        }).forEach(word -> {
            Utils.sleep(50);
            if (resultMap.containsKey(word)) {
                resultMap.put(word, resultMap.get(word) + 1L);
            } else {
                resultMap.put(word, 1L);
            }
        });
        long end = System.nanoTime();
        System.out.printf("time spent for wordcount by parallel stream: %.3f sec.%n", (end - start) * 1E-9);
        resultMap.forEach(MxWordCountByParalleStreamTest::accept);
    }

    private static void accept(String word, Long count) {
        System.out.println(word + ": " + count);
    }

执行结果:

time spent for wordcount by parallel stream: 1.999 sec.
over: 2
a: 2
away: 2
nature: 2
jumped: 2
i: 2
seven: 2
cow: 2
am: 2
an: 2
two: 2
dwarfs: 2
the: 8
doctor: 2
apple: 2
with: 2
moon: 2
at: 2
white: 2
snow: 2
and: 2
keep: 2
day: 2
time spent for wordcount by sequence stream: 3.850 sec.
over: 2
a: 2
away: 2
nature: 2
jumped: 2
seven: 2
i: 2
cow: 2
am: 2
an: 2
two: 2
dwarfs: 2
the: 8
doctor: 2
apple: 2
with: 2
moon: 2
at: 2
white: 2
snow: 2
and: 2
keep: 2
day: 2

Process finished with exit code 0

可见使用并行流可以节省一半的时间.

需要特别注意的是, 不是任何情况下使用并行流都可以节省时间, 并且使用时还要特别当心是否有线程安全问题.

Views: 560

Kafka的消息存储

为了便于说明问题,假设这里只有一个单节点的伪分布式Kafka集群。在这个Kafka broker实例的$KAFKA_HOME/config/server.properties中配置 log.dirs=/tmp/kafka-logs,以此来设置 Kafka 消息文件存储目录。并通过命令:

$KAFKA_HOME/bin/kafka-topics.sh --create --zookeeper localhost:2181 --partitions 4 --topic topic_test --replication-factor 1

创建一个 topic:topic_test,partition的数量配置为4。接下来可以在 /tmp/kafka-logs 目录中可以看到生成了 4 个partition目录:

drwxr-xr-x 2 root root 4096 Apr 10 16:10 topic_test-0
drwxr-xr-x 2 root root 4096 Apr 10 16:10 topic_test-1
drwxr-xr-x 2 root root 4096 Apr 10 16:10 topic_test-2
drwxr-xr-x 2 root root 4096 Apr 10 16:10 topic_test-3

在Kafka文件存储中,同一个topic下有多个不同的partition,每个partiton为一个目录,partition的名称规则为topic名称+有序序号,第一个序号从0开始计,最大的序号为partition数量减1,类似数组的索引, partition是实际物理上的概念,而topic是逻辑上的概念,更多表象是一个消息的类别,当然,partition还可以细分为segment(文件),一个partition物理上由多个segment组成.

为什么不能以partition 作为存储单位?

如果就以partition 为最小存储单位,可以想象,当Kafka producer不断发送消息,必然会引起partition文件的无限扩张,将对消息文件的维护以及已消费的消息的清理带来严重的影响,新数据是在文件尾部追加的,不论文件数据文件有多大,这个操作永远都是 O(1)的查找,再者,查找某个offset的Message是顺序查找的。因此,如果数据文件很大的话,查找的效率就低. 因此,需以segment 为单位将partition进一步细分。每个partition相当于一个巨型文件被平均分配到多个大小相等的segment数据文件中(每个segment文件中消息数量不一定相等)这种特性也方便old segment的删除,即方便已被消费的消息的清理,提高磁盘的利用率。每个partition只需要支持顺序读写就行.segment的文件生命周期由服务端配置参数log.segment.byteslog.roll.{ms,hours} 等相关参数决定。

segment文件由两部分组成,分别为 .index文件和.log文件,分别表示为segment索引文件和数据文件。这两个文件的命令规则为:partition全局的第一个segment从0开始,后续每个segment文件名为上一个segment文件最后一条消息的offset值,数值大小为64位,20位数字字符长度,没有数字用0填充,如下:

00000000000000000000.index
00000000000000000000.log
00000000000000000000.timeindex
00000000000000170410.index
00000000000000170410.log
00000000000000170410.timeindex
00000000000000239430.index
00000000000000239430.log
00000000000000239430.timeindex

以上面的segment文件为例,展示出segment中的00000000000000170410.index文件和00000000000000170410.log文件是存在对应关系的,.index索引文件存储大量的元数据,.log数据文件存储大量的消息,索引文件中的元数据指向对应数据文件中message的物理偏移地址(也就是实际的偏移地址,因为会涉及到segment文件清理)。其中以.index索引文件中的元数据[3, 348]为例,在.log数据文件表示第3个消息,该消息的物理偏移地址为348。

如何保证消息消费的有序性呢?

举个例子,比如说生产者生产了25个订单,订单假设分为创建-提交-付款-发货4个步骤,那么消费者在消费的时候按照0到100这个从小到大的顺序消费,那么kafka如何保证这种有序性呢?

难度就在于,生产者生产出0到100这100条数据之后,通过一定的分组策略存储到broker的partition中的时候,比如0到10这10条消息被存到了这个partition中,10到20这10条消息被存到了那个partition中,这样的话,消息在分组存到partition中的时候就已经被分组策略搞得无序了。

那么能否做到消费者在消费消息的时候全局有序呢?遇到这个问题,我们可以回答,在大多数情况下是做不到全局有序的。但在某些情况下是可以做到的。比如我的partition只有一个,这种情况下是可以全局有序的。

那么可能有人又有疑问了,只有一个partition的话,哪里来的分布式呢?哪里来的负载均衡呢?所以说,全局有序是一个伪命题!让订单全局有序根本没有办法在kafka内部做到。但是我们只能保证当前这个partition内部消息消费的有序性。

针对一个topic里面的数据,只能做到partition内部有序,不能做到全局有序。

当然不使用一个partition我们可以从代码层面解决.

如何从partition中通过offset查找message?

以上图为例,读取offset=170418的消息,首先查找segment文件,其中00000000000000000000.index为最开始的文件,第二个文件为00000000000000170410.index(起始偏移为170410),而第三个文件为00000000000000239430.index(起始偏移为239430),所以这个offset=170418就落到了第二个文件之中。其他后续文件可以依次类推,以其实偏移量命名并排列这些文件,然后根据二分查找法就可以快速定位到具体文件位置。其次根据00000000000000170410.index文件中的[8,1325]定位到00000000000000170410.log文件中的1325的位置进行读取。

怎么知道何时读完本条消息,否则就读到下一条消息的内容了?
这个问题就得引出kafka的消息结构,如下图所示:

file

消息都具有固定的物理结构,包括:offset(8 Bytes)、消息体的大小(4 Bytes)、crc32(4 Bytes)、magic(1 Byte)、attributes(1 Byte)、key length(4 Bytes)、key(K Bytes)、payload(N Bytes)等等字段,可以确定一条消息的大小,即读取到哪里截止。

总结,offset的查找方机制是建立在offset是有序的,索引文件被映射到内存中,所以查找的速度还是很快的。

另外,Kafka的Message存储采用了分区(partition),Segment和index这几个手段来达到了高效性。

Views: 602

Kafka如何保证无消息丢失

配置Kafka无消息丢失

Kafka 只对“已提交”的消息(committed message)做有限度的持久化保证。Kafka 的 一个Broker或多个Broker 成功地接收到一条消息并写入到日志文件后,它们会告诉生产者程序这条消息已成功提交,具体是一个Broker还是多个Broker取决ack参数的配置。

要想要消息不丢失,假如你的消息保存在 N 个 Kafka Broker 上,那么这个前提条件就是这 N 个 Broker 中至少有 1 个存活。

目前 Kafka Producer 的send方法是异步发送消息的,也就是说如果你调用的是 producer.send(msg) 这个 API,那么它通常会立即返回,消息会异步继续发送到Broker中,在这个过程中是有可能失败的,比如网络抖动导致发送失败,比如消息格式不正确等导致Broker拒绝消息。

解决方法:Producer 要使用带有回调通知的发送 API,也就是说不要使用 producer.send(msg),而要使用 producer.send(msg, callback)。

通过callback(回调),它能准确地告诉你消息是否真的提交成功了。一旦出现消息提交失败的情况,你就可以有针对性地进行处理。举例来说,如果是因为那些瞬时错误,那么仅仅让 Producer 重试就可以了;如果是消息不合格造成的,那么可以调整消息格式后再次发送。总之,处理发送失败的责任在 Producer 端而非 Broker 端。

import org.apache.kafka.clients.producer.*;
import org.apache.kafka.common.serialization.StringSerializer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.util.Properties;
import java.util.concurrent.ExecutionException;

public class ProducerWithCallBackExample {

    public static void main(String[] args) throws ExecutionException, InterruptedException {

        final Logger logger = LoggerFactory.getLogger(ProducerWithCallBackExample.class);

        //Create producer properties
        Properties properties = new Properties();
        properties.setProperty(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "hadoop000:9092");
        properties.setProperty(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        properties.setProperty(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());

        //Create the Producer
        KafkaProducer<String, String> producer = new KafkaProducer<String, String>(properties);

        //Send 10 messages
        for (int i = 0; i < 10; i++) {
            //Create a Producer Record
            ProducerRecord<String, String> record = new ProducerRecord<String, String>("my-topic", "id_" + i, "Hello Kakfka" + i);

            logger.info("Key :" + "id_" + i); //Log the Key
            //send data
            producer.send(record, new Callback() {
                @Override
                public void onCompletion(RecordMetadata recordMetadata, Exception e) {
                    if (e == null) {
                        logger.info("Received new metadata. \n" +
                                "Topic:" + recordMetadata.topic() + "\n" +
                                "Pratition:" + recordMetadata.partition() + "\n" +
                                "Offeset:" + recordMetadata.offset() + "\n" +
                                "Timestamp" + recordMetadata.timestamp());
                    } else {
                        logger.error("Error while producing", e);
                    }
                }
            });
           producer.flush();
        }
        producer.close();
    }
}

最佳实践

分享一下 Kafka 无消息丢失的配置的最佳实践:

  1. 不要使用 producer.send(msg),而要使用 producer.send(msg, callback)。记住,一定要使用带有回调通知的 send 方法。
  2. 设置 acks = all。acks 是 Producer 的一个参数,代表了你对“已提交”消息的定义。如果设置成 all,则表明所有副本 Broker 都要接收到消息,该消息才算是“已提交”。这是最高等级的“已提交”定义。
  3. 设置 retries 为一个较大的值。这里的 retries 同样是 Producer 的参数,对应前面提到的 Producer 自动重试。当出现网络的瞬时抖动时,消息发送可能会失败,此时配置了 retries > 0 的 Producer 能够自动重试消息发送,避免消息丢失。
  4. 设置 unclean.leader.election.enable = false。这是 Broker 端的参数,它控制的是哪些 Broker 有资格竞选分区的 Leader。如果一个 Broker 落后原先的 Leader 太多,那么它一旦成为新的 Leader,必然会造成消息的丢失。故一般都要将该参数设置成 false,即不允许这种情况的发生。
  5. 设置 replication.factor >= 3。这也是 Broker 端的参数。其实这里想表述的是,最好将消息多保存几份,毕竟目前防止消息丢失的主要机制就是冗余。
  6. 设置 min.insync.replicas > 1。这依然是 Broker 端参数,控制的是消息至少要被写入到多少个副本才算是“已提交”。设置成大于 1 可以提升消息持久性。在实际环境中千万不要使用默认值 1。
  7. 确保 replication.factor > min.insync.replicas。如果两者相等,那么只要有一个副本挂机,整个分区就无法正常工作了。我们不仅要改善消息的持久性,防止数据丢失,还要在不降低可用性的基础上完成。推荐设置成 replication.factor = min.insync.replicas + 1
  8. 确保消息消费完成再提交。Consumer 端有个参数 enable.auto.commit,最好把它设置成 false,并采用手动提交位移的方式。就像前面说的,这对于单 Consumer 多线程处理的场景而言是至关重要的。

Views: 303

Storm DRCP应用(计算推特Reach值)

需求

针对twitter网站上的一篇推文的接触用户(也叫REACH值)进行统计。

Reach值让你了解推文的真实覆盖到的用户群体, 要计算一个推文URL的Reach值,需要以下4步:

  1. 根据推文的URL查询数据库获取全部直接接触用户(转发的用户)
  2. 再根据接触用户通过查询数据库获取每个用户的全部粉丝
  3. 对粉丝集合中的用户进行去重处理
  4. 最后统计去重后的用户数, 即这个推文的Reach值

拓扑定义

一个单独的Reach计算在计算期间可能涉及到数千次数据库访问和数千万的粉丝记录查询,可能是一个非常耗时的计算。在storm上实现这个功能非常简单。在一台机器上,Reach计算可能花费数分钟,而在storm集群,最难计算Reach的URL也只需数秒。

Storm-starter 项目中有一个计算Reach样例,Reach拓扑定义如下所示:

LinearDRPCTopologyBuilder builder = new LinearDRPCTopologyBuilder("reach"); 

builder.addBolt(new GetTweeters(), 3); 
builder.addBolt(new GetFollowers(), 2) .shuffleGrouping(); 
builder.addBolt(new PartialUniquer(), 3) .fieldsGrouping(new Fields("id", "follower")); 
builder.addBolt(new CountAggregator()) .fieldsGrouping(new Fields("id")); 

过程分析

以计算https://tech.backtype.com/blog/123的Reach值为例进行说明:

image-20211024202541029

这个拓扑以4个步骤的形式执行, 因此一共有4个Bolt:

  1. GetTweeters 从数据库获取给定URL对应的用户并发射出去
  2. GetFollowers 从数据库获取每个用户对应的粉丝并发射出去
  3. PartialUniquer按粉丝进行字段分组并利用Set集合进行去重处理, 将去重后的粉丝数量进行局部累加, 并把结果发送出去
  4. 最后,CountAggregator从每个PartialUniquer任务接收计数并对累加求和作为返回值给DRPC客户端。

项目代码

为了方便项目中数据库使用Map集合进行伪造, 代码如下:

package storm.example.drpc;

/**
 * 计算推文的REACH值
 * storm.example.drpc.AdvanceDRPCTopology
 */

import org.apache.storm.Config;
import org.apache.storm.LocalCluster;
import org.apache.storm.LocalDRPC;
import org.apache.storm.StormSubmitter;
import org.apache.storm.coordination.BatchOutputCollector;
import org.apache.storm.drpc.LinearDRPCTopologyBuilder;
import org.apache.storm.task.TopologyContext;
import org.apache.storm.topology.BasicOutputCollector;
import org.apache.storm.topology.OutputFieldsDeclarer;
import org.apache.storm.topology.base.BaseBasicBolt;
import org.apache.storm.topology.base.BaseBatchBolt;
import org.apache.storm.tuple.Fields;
import org.apache.storm.tuple.Tuple;
import org.apache.storm.tuple.Values;
import org.apache.storm.utils.Utils;

import java.util.*;

// 计算一篇推文的REACH值
public class AdvanceDRPCTopology {

    // MAP: URL -> USER
    public static Map<String, List<String>> TWEETERS_DB = new HashMap<String, List<String>>() {{
        put("foo.com/blog/1", Arrays.asList("sally", "bob", "tim", "george", "nathan"));
        put("engineering.twitter.com/blog/5", Arrays.asList("adam", "david", "sally", "nathan"));
        put("tech.backtype.com/blog/123", Arrays.asList("tim", "mike", "john"));
    }};

    // MAP: USER -> FOLLOWERS
    public static Map<String, List<String>> FOLLOWERS_DB = new HashMap<String, List<String>>() {{
        put("sally", Arrays.asList("bob", "tim", "alice", "adam", "jim", "chris", "jai"));
        put("bob", Arrays.asList("sally", "nathan", "jim", "mary", "david", "vivian"));
        put("tim", Arrays.asList("alex"));
        put("nathan", Arrays.asList("sally", "bob", "adam", "harry", "chris", "vivian", "emily", "jordan"));
        put("adam", Arrays.asList("david", "carissa"));
        put("mike", Arrays.asList("john", "bob"));
        put("john", Arrays.asList("alice", "nathan", "jim", "mike", "bob"));
    }};

    public static LinearDRPCTopologyBuilder construct() {

        LinearDRPCTopologyBuilder builder = new LinearDRPCTopologyBuilder("reach");

        builder.addBolt(new GetTweeters(), 3);
        builder.addBolt(new GetFollowers(), 2).shuffleGrouping();
        builder.addBolt(new PartialUniquer(), 3).fieldsGrouping(new Fields("id", "follower"));
        builder.addBolt(new CountAggregator(), 1).fieldsGrouping(new Fields("id"));

        return builder;
    }

    public static void main(String[] args) throws Exception {

        LinearDRPCTopologyBuilder builder = construct();
        Config conf = new Config();

        if (args == null || args.length == 0) {
            // conf.setMaxTaskParallelism(3);
            LocalDRPC drpc = new LocalDRPC();
            LocalCluster cluster = new LocalCluster();
            cluster.submitTopology("reach-drpc", conf, builder.createLocalTopology(drpc));

            String[] urlsToTry = new String[]{"foo.com/blog/1", "engineering.twitter.com/blog/5", "notaurl.com", "tech.backtype.com/blog/123"};

            for (String url : urlsToTry) {
                System.err.println("Thread[" + Thread.currentThread().getName() + "] " +"Reach of " + url + ": " + drpc.execute("reach", url));
                Utils.sleep(10000);
            }
            cluster.shutdown();
            drpc.shutdown();
        } else {
            conf.setNumWorkers(1); //6
            StormSubmitter.submitTopology(args[0], conf, builder.createRemoteTopology());
        }
    }

    // BOLT 1: 根据URL发送推特用户到数据流中
    public static class GetTweeters extends BaseBasicBolt {

        @Override
        public void execute(Tuple tuple, BasicOutputCollector collector) {

            Object id = tuple.getValue(0);
            String url = tuple.getString(1);
            List<String> tweeters = TWEETERS_DB.get(url);

            if (tweeters != null) {
                for (String tweeter : tweeters) {
                    collector.emit(new Values(id, tweeter));
                }
            }
        }

        @Override
        public void declareOutputFields(OutputFieldsDeclarer declarer) {

            declarer.declare(new Fields("id", "tweeter"));
        }

    }

    // BOLT 2: 根据推特用户发送对应的粉丝到数据流中
    public static class GetFollowers extends BaseBasicBolt {

        @Override
        public void execute(Tuple tuple, BasicOutputCollector collector) {

            Object id = tuple.getValue(0);
            String tweeter = tuple.getString(1);
            List<String> followers = FOLLOWERS_DB.get(tweeter);

            if (followers != null) {
                for (String follower : followers) {
                    System.err.println("Thread[" + Thread.currentThread().getName() + "] " + "request-id[" + id + "]: " + tweeter + "'s follower: " + follower);
                    collector.emit(new Values(id, follower));
                }
            }
        }

        @Override
        public void declareOutputFields(OutputFieldsDeclarer declarer) {

            declarer.declare(new Fields("id", "follower"));
        }
    }

    // BOLT 3: 对粉丝去重后的计数发送的数据流中
    public static class PartialUniquer extends BaseBatchBolt {

        BatchOutputCollector collector;
        Object id;
        Set<String> followers = new HashSet<String>();

        @Override
        public void prepare(Map conf, TopologyContext context, BatchOutputCollector collector, Object id) {
            this.collector = collector;
            this.id = id;
        }

        @Override
        public void execute(Tuple tuple) {
            followers.add(tuple.getString(1));
        }

        @Override
        public void finishBatch() {
            System.err.println("Thread[" + Thread.currentThread().getName() + "] " +"request-id[" + id + "]: " + "followers.size(): " + followers.size() + followers);
            collector.emit(new Values(id, followers.size()));
        }

        @Override
        public void declareOutputFields(OutputFieldsDeclarer declarer) {
            declarer.declare(new Fields("id", "partial-count"));
        }
    }

    // BOLT 4:将所有计数累加的结果发送的数据流中
    public static class CountAggregator extends BaseBatchBolt {

        BatchOutputCollector collector;

        Object id;
        int reach = 0;

        @Override
        public void prepare(Map conf, TopologyContext context, BatchOutputCollector collector, Object id) {
            this.collector = collector;
            this.id = id;
        }

        @Override
        public void execute(Tuple tuple) {
            int partiaLReach = tuple.getInteger(1);
            System.err.println("Thread[" + Thread.currentThread().getName() + "] " +"request-id[" + id + "]: " + "Partial reach is " + partiaLReach);
            reach += tuple.getInteger(1);
        }

        @Override
        public void finishBatch() {
            System.err.println("Thread[" + Thread.currentThread().getName() + "] " +"request-id[" + id + "]: " + "Global reach is " + reach);
            collector.emit(new Values(id, reach));
        }

        @Override
        public void declareOutputFields(OutputFieldsDeclarer declarer) {
            declarer.declare(new Fields("id", "reach"));
        }
    }
}

Reach计算的每一单步都是可以并行执行的,而且定义一个DRPC拓扑也非常简单。

需要注意的是, CountAggregator 只有在 PartialUniquer 并行度大于1个时候才有意义.

另外CountAggregatorPartialUniquer 这两个Bolt都是继承的BaseBatchBolt,这种Bolt发射数据流的时机是在相同id的数据都处理完成后, 才会触发(类似事务).

谢谢!

Views: 356

Kafka Producer API

http://kafka.apache.org/25/documentation.html#api

Kafka 架构回顾

image-20211011145945791

1)Producer :消息生产者,就是向kafka broker发消息的客户端;

2)Consumer :消息消费者,向kafka broker取消息的客户端;

3)Topic :可以理解为一个队列;

4) Consumer Group (CG):这是kafka用来实现一个topic消息的广播(发给所有的consumer)和单播(发给任意一个consumer)的手段。一个topic可以有多个CG。topic的消息会复制(不是真的复制,是概念上的)到所有的CG,但每个partion只会把消息发给该CG中的一个consumer。如果需要实现广播,只要每个consumer有一个独立的CG就可以了。要实现单播只要所有的consumer在同一个CG。用CG还可以将consumer进行自由的分组而不需要多次发送消息到不同的topic;

5)Broker :一台kafka服务器就是一个broker。一个集群由多个broker组成。一个broker可以容纳多个topic;

6)Partition:为了实现扩展性,一个非常大的topic可以分布到多个broker(即服务器)上,一个topic可以分为多个partition,每个partition是一个有序的队列。partition中的每条消息都会被分配一个有序的id(offset)。kafka只保证按一个partition中的顺序将消息发给consumer,不保证一个topic的整体(多个partition间)的顺序;

7)Offset:kafka的存储文件都是按照offset.kafka来命名,用offset做名字的好处是方便查找。例如你想找位于2049的位置,只要找到2048.kafka的文件即可。当然the first offset就是00000000000.kafka。

添加 maven 依赖

Producer API可以让应用程序发送数据流到Kafka集群的消息主题中。

Producer的使用可以查看官方文档 javadocs.

http://kafka.apache.org/25/javadoc/index.html?org/apache/kafka/clients/producer/KafkaProducer.html

要使用Producer可以添加以下Maven依赖:

<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-clients</artifactId>
    <version>2.5.0</version>
</dependency>

Producer的简单使用

public class KafkaProducer<K, V> implements Producer<K, V>

KafkaProducer是Kafka客户端,用于将一条消息(记录)发布到Kafka集群。

KafkaProducer是线程安全的,跨线程共享单个Producer实例通常比拥有多个Producer实例更快。

下面是一个使用Producer发送消息的简单示例,该记录使用包含0到99的数字序列的字符串作为键/值对。

 Properties props = new Properties();
 props.put("bootstrap.servers", "localhost:9092");
 props.put("acks", "all");
 props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
 props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");

 Producer<String, String> producer = new KafkaProducer<>(props);
 for (int i = 0; i < 100; i++)
     producer.send(new ProducerRecord<String, String>("my-topic", Integer.toString(i), Integer.toString(i)));

 producer.close();

注意这里send()发送消息是异步的.

Producer 生产过程

写入方式

producer采用推(push)模式将消息发布到broker,每条消息都被追加(append)到分区(patition)中,属于顺序写磁盘(顺序写磁盘效率比随机写内存要高,保障kafka吞吐率)。

分区(Partition)

消息发送时都被发送到一个topic,其本质就是一个目录,而topic是由一些Partition Logs(分区日志)组成,其组织结构如下图所示:

image-20211011150538034

我们可以看到,每个Partition中的消息都是有序的,生产的消息被不断追加到Partition log上,其中的每一个消息都被赋予了一个唯一的offset值。

如何唯一标识一个消息: 主题->分区->偏移量

1)分区的原因

(1)方便在集群中扩展,每个Partition可以通过调整以适应它所在的机器,而一个topic又可以有多个Partition组成,因此整个集群就可以适应任意大小的数据了;

(2)可以提高并发,因为可以以Partition为单位读写了。

2)分区的原则

(1)指定了patition,则直接使用;

(2)未指定patition但指定key,通过对key的value进行hash出一个patition;

(3)patition和key都未指定,使用轮询选出一个patition。

DefaultPartitioner class {

public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) {

    List<PartitionInfo> partitions = cluster.partitionsForTopic(topic);

    int numPartitions = partitions.size();

    if (keyBytes == null) {

      int nextValue = nextValue(topic);

      List<PartitionInfo> availablePartitions = cluster.availablePartitionsForTopic(topic);

      if (availablePartitions.size() > 0) {

        int part = Utils.toPositive(nextValue) % availablePartitions.size();

        return availablePartitions.get(part).partition();

      } else {

        // no partitions are available, give a non-available partition

        return Utils.toPositive(nextValue) % numPartitions;

      }

    } else {

      // hash the keyBytes to choose a partition

      return Utils.toPositive(Utils.murmur2(keyBytes)) % numPartitions;

    }

  }

}

副本(Replication)

同一个partition可能会有多个replication(对应 server.properties 配置中的 default.replication.factor=N)。没有replication的情况下,一旦broker 宕机,其上所有 patition 的数据都不可被消费,同时producer也不能再将数据存于其上的patition。引入replication之后,同一个partition可能会有多个replication,而这时需要在这些replication之间选出一个leader,producer和consumer只与这个leader交互,其它replication作为follower从leader 中复制数据。

消息写入流程

image-20211011150947554

producer写入消息流程如下:

1)producer先从zookeeper的 "/brokers/.../state"节点找到该partition的leader

2)producer将消息发送给该leader

3)leader将消息写入本地log

4)followers从leader pull消息,写入本地log后向leader发送ACK

5)leader收到所有ISR中的replication的ACK后,增加HW(high watermark,最后commit 的offset)并向producer发送ACK

Broker 保存消息

存储方式

物理上把topic分成一个或多个patition(对应 server.properties 中的num.partitions=3配置),每个patition物理上对应一个文件夹(该文件夹存储该patition的所有消息和索引文件),如下:

[hadoop@hadoop100 logs]$ ll

drwxrwxr-x. 2 hadoop hadoop 4096 8月  6 14:37 first-0
drwxrwxr-x. 2 hadoop hadoop 4096 8月  6 14:35 first-1
drwxrwxr-x. 2 hadoop hadoop 4096 8月  6 14:37 first-2

[hadoop@hadoop100 logs]$ cd first-0

[hadoop@hadoop100 first-0]$ ll

-rw-rw-r--. 1 hadoop hadoop 10485760 8月  6 14:33 00000000000000000000.index
-rw-rw-r--. 1 hadoop hadoop   219 8月  6 15:07 00000000000000000000.log
-rw-rw-r--. 1 hadoop hadoop 10485756 8月  6 14:33 00000000000000000000.timeindex
-rw-rw-r--. 1 hadoop hadoop    8 8月  6 14:37 leader-epoch-checkpoint

存储策略

无论消息是否被消费,kafka都会保留所有消息。有两种策略可以删除旧数据:

1)基于时间:log.retention.hours=168

2)基于大小:log.retention.bytes=1073741824

需要注意的是,因为Kafka读取特定消息的时间复杂度为O(1),即与文件大小无关,所以这里删除过期文件与提高 Kafka 性能无关。

Zookeeper存储结构

image-20211018163259464

注意:producer不在zk中注册,消费者在zk中注册。

producer 包括一个保存尚未传输到服务器的消息(record)的缓冲区空间池,以及一个负责将这些记录转换为请求并将其传输到集群的后台I/O线程。producer使用后未及时关闭将导致这些资源泄漏。

send()方法是异步的,调用send()方法会将待发送消息(record)添加到缓冲区等待发送,并立即返回。生产者会积攒一批消息后批量发送以提高效率。

acks配置项用来控制请求如何被认为已处理完成。我们指定的“all”设置将导致阻塞全部记录提交,这是最慢但最持久的设置。

如果请求失败,生产者可以自动重试,除非配置了retries配置项为0。启用重试可能导致出现重复消息的可能性(有关消息传递语义的详细信息,请参阅文档)。

Producer 在每个分区维都维护未发送记录的缓冲区。这些缓冲区的大小由batch.size配置项进行配置,增大这个配置的值可以允许一个批次处理更多消息,但需要更多的内存(因为我们通常会为每个活动分区拥有一个这样的缓冲区)。

默认情况下,即使缓冲区中有额外的未使用空间,缓冲数据也是可以立即发送出去的。然而如果你想减少请求的数量可以设置linger.ms为一个大于0的数, 这可以让Producer发送请求之前等待该毫秒数,这样可以有更多消息在一个批次中批量发送。

在上面的代码片段中,可能会在一个请求中发送所有100条记录,因为我们将延迟时间设置为1毫秒。但是,如果我们没有填满缓冲区,这个设置会给请求增加1毫秒的延迟,以等待更多的记录填充到缓冲区。因此在负载大的情况下,无论是否配置linger.ms,消息都是通过批处理发送的;将linger.ms的值设置为大于0的值的意义在于,在没有处于最大负载的情况下可以有效减少请求数量, 代价是增加一点点延迟。

配置项buffer.memory 用来控制Producer用于缓冲区的可用内存总量。如果消息发送到服务器的速度大于缓冲区耗尽的速度,当缓冲区空间耗尽时,其他消息的发送将被阻塞,阻塞时间的阈值由max.block决定。之后,它会抛出一个TimeoutException异常。

配置项 key.serializervalue.serializer 用来配置键值对象的序列化类, 对于简单的字符串或字节类型,可以使用Kafka自带的ByteArraySerializer或StringSerializer。

从Kafka 0.11这个版本开始,KafkaProducer支持另外两种模式:幂等Producer和事务Producer。幂等Producer加强了卡夫卡的交付语义,从至少一次交付(at least once delivery)加强为准确地一次交付(exactly once delivery), 幂等Producer的消息无论失败重试多少次都不会导致数据重复统计。而事务Producer是以原子操作将消息发送到多个分区(和主题!),即这些操作要么同时成功, 要么同时失败!.

若要启用等幂性(idempotence),则 enable.idempotence配置必须设置为true。此时retries配置值将默认为Integer.MAX_VALUE, 而acks配置默认为all。使用幂等性Producer在API上没有变化,因此不需要修改现有应用程序。

一旦启用了幂等性(idempotence),建议不要配置重试(retries)参数,因为它默认值为Integer.MAX_VALUE。另外,如果send(ProducerRecord)方法在进行无限次重试之后仍然返回错误(例如消息在缓冲区中过期),那么建议关闭Producer并检查最后生成的消息的内容,以确保没有重复消息。最后,生产者只能保证在单个会话中发送的消息的幂等性。

过时的Producer API

一些公司可能因为历史原因, 仍然在使用旧的(过时的)API.

pom.xml

<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka_2.11</artifactId>
    <version>0.10.0.1</version>
</dependency>

简单代码示例

BasicConfigurator.configure();
Properties properties = new Properties();
properties.put("metadata.broker.list", "hadoop000:9092");
properties.put("request.required.acks", "1");
properties.put("serializer.class", "kafka.serializer.StringEncoder");

Producer<Integer, String> producer = new Producer<>(new ProducerConfig(properties));

KeyedMessage<Integer, String> message = new KeyedMessage<>("t1", "hello world");
producer.send(message);

过时的生产者API不支持幂等性和事务.

复杂代码示例

自定义调度器代码

package com.niit.kafka.producer.newapi;

/**
 * @Author: deLucia
 * @Date: 2021/10/7
 * @Version: 1.0
 * @Description:
 * 使用更多配置项的生产者
 * 启动自定义分区: 不管有多少分区,只写到分区0
 * 先开启消费者, 运行程序, 并测试 tail -F my-topic-3-0/00000000000000000000.log
 */

import org.apache.kafka.clients.producer.Partitioner;
import org.apache.kafka.common.Cluster;

import java.util.Map;

public class CustomPartitioner implements Partitioner {

    @Override
    public void configure(Map<String, ?> configs) {

    }

    @Override
    public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) {
        // 控制分区
        return 0;
    }

    @Override
    public void close() {

    }
}

生产者代码


/**
 * @Author: deLucia
 * @Date: 2021/10/6
 * @Version: 1.0
 * @Description: old api from kafka_2.11
 */

import kafka.javaapi.producer.Producer;
import kafka.producer.KeyedMessage;
import kafka.producer.ProducerConfig;
import org.apache.kafka.common.utils.Utils;
import org.apache.log4j.BasicConfigurator;

import java.util.Properties;

public class DetailedProducerExample {

    @SuppressWarnings("deprecation")
    public static void main(String[] args) {
        BasicConfigurator.configure();
        Properties props = new Properties();
        // Kafka服务端的主机名和端口号
        props.put("metadata.broker.list", "hadoop100:9092");
        // 等待所有副本节点的应答
        props.put("request.required.acks", "0");
        // 序列化类
        props.put("serializer.class", "kafka.serializer.StringEncoder");
        // partitioner
        //      - default: kafka.producer.DefaultPartitioner
        //          - hash(key)%partitionNum)
        //      - custom partitioner:
        //          - send to partition #0
        props.put("partitioner.class", CustomPartitioner.class.getName());

        Producer<String, String> producer = new Producer<>(new ProducerConfig(props));

        for (int i = 0; i < 10; i++) {
            KeyedMessage<String, String> message = new KeyedMessage<>("my-topic-3", "hello world " + i);
            producer.send(message);
            Utils.sleep(1000);
        }
        producer.close();
    }
}

Producer 配置参数

所有producer的配置参阅官网

同步发送结合自定义调度器

自定义调度器

package com.niit.kafka.producer.newapi;

/**
 * @Author: deLucia
 * @Date: 2021/10/7
 * @Version: 1.0
 * @Description:
 * 使用更多配置项的生产者
 * 启动自定义分区: 不管有多少分区,只写到分区0
 * 先开启消费者, 运行程序, 并测试 tail -F my-topic-3-0/00000000000000000000.log
 */

import org.apache.kafka.clients.producer.Partitioner;
import org.apache.kafka.common.Cluster;

import java.util.Map;

public class CustomPartitioner implements Partitioner {

    @Override
    public void configure(Map<String, ?> configs) {

    }

    @Override
    public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) {
        // 控制分区
        return 0;
    }

    @Override
    public void close() {

    }
}

生产者代码

package com.niit.kafka.producer.newapi;

/**
 * @Author: deLucia
 * @Date: 2020/8/22
 * @Version: 1.0
 * @Description: new Producer APi from kafka_client
 * STEPS:
    create topic t1 with multiple partitioners.
    bin/kafka-console-consumer.sh --bootstrap-server hadoop000:9092,hadoop000:9093,hadoop000:9094 --topic t1
    $ tail -F t1-0/0000000000000000000X.log
 */

import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.clients.producer.RecordMetadata;
import org.apache.kafka.common.serialization.StringSerializer;
import org.apache.log4j.BasicConfigurator;
import org.apache.kafka.clients.producer.internals.DefaultPartitioner;

import java.util.Properties;
import java.util.concurrent.ExecutionException;

public class DetailedProducerExample {
    public static void main(String[] args) throws ExecutionException, InterruptedException {
        BasicConfigurator.configure();

        // Create producer properties
        Properties properties = new Properties();
        //      - kafka clustger
        properties.setProperty(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "hadoop000:9092,hadoop000:9093,hadoop000:9094");
        //      - ack level
        //              no wait = 0
        //              wait leader  = 1
        //              wait all     = -1
        properties.setProperty(ProducerConfig.ACKS_CONFIG, "all");
        //      - retries
        properties.setProperty(ProducerConfig.RETRIES_CONFIG, "0");
        //      - batch.size
        properties.setProperty(ProducerConfig.BATCH_SIZE_CONFIG, "16384");
        //      - linger.ms
        properties.setProperty(ProducerConfig.LINGER_MS_CONFIG, "1");
        //      - buffer.memory
        properties.setProperty(ProducerConfig.BUFFER_MEMORY_CONFIG, "33554421");
        //      - KV serializer classes
        properties.setProperty(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        properties.setProperty(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        //      - Defauilt Partitioner
        //           default: org.apache.kafka.clients.producer.internals.DefaultPartitioner;
        //              - 如果指定分区则使用指定的分区
        //              - 如果没有指定分区但是指定key, 则计算哈希取模得到分区索引
        //              - 如果没有指定分区也没有指定key, 则采用轮询算法
        //      - Custom Partitioner
        //              - 总是发送到分区0上
        properties.setProperty(ProducerConfig.PARTITIONER_CLASS_CONFIG, CustomPartitioner.class.getName());
        // Create the Producer
        KafkaProducer<String, String> producer = new KafkaProducer<>(properties);

        // send data
        for (int i = 0; i < 10; i++) {
            // Create a Producer Record
            ProducerRecord<String, String> record = new ProducerRecord<>("my-topic-3", String.valueOf(i));
            RecordMetadata recordMetadata = producer.send(record).get();// 同步发送
            int partition = recordMetadata.partition();
            long offset = recordMetadata.offset();
            System.out.println("Partition[" + partition + "], offset[" + offset + "], message[" + i + "]");

// producer.send(record); // 异步发送
            producer.flush();
        }

        producer.close();
    }
}

异步发送结合回调

这里回调采用使用匿名内部类的方式:

import org.apache.kafka.clients.producer.*;
import org.apache.kafka.common.serialization.StringSerializer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.util.Properties;
import java.util.concurrent.ExecutionException;

public class ProducerWithCallBackExample {

    public static void main(String[] args) throws ExecutionException, InterruptedException {

        final Logger logger = LoggerFactory.getLogger(ProducerWithCallBackExample.class);

        //Create producer properties
        Properties properties = new Properties();
        properties.setProperty(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "hadoop000:9092");
        properties.setProperty(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        properties.setProperty(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());

        //Create the Producer
        KafkaProducer<String, String> producer = new KafkaProducer<String, String>(properties);

        //Send 10 messages
        for (int i = 0; i < 10; i++) {
            //Create a Producer Record
            ProducerRecord<String, String> record = new ProducerRecord<String, String>("my-topic", "id_" + i, "Hello Kakfka" + i);

            logger.info("Key :" + "id_" + i); //Log the Key
            //send data
            producer.send(record, new Callback() {
                @Override
                public void onCompletion(RecordMetadata recordMetadata, Exception e) {
                    if (e == null) {
                        logger.info("Received new metadata. \n" +
                                "Topic:" + recordMetadata.topic() + "\n" +
                                "Pratition:" + recordMetadata.partition() + "\n" +
                                "Offeset:" + recordMetadata.offset() + "\n" +
                                "Timestamp" + recordMetadata.timestamp());
                    } else {
                        logger.error("Error while producing", e);
                    }
                }
            });
           producer.flush();
        }
        producer.close();
    }
}

事务Producer

客户端能与之交互的Broker上安装Kafka版本要至少0.10.0。有些代理服务器可能不支持某些客户端特性。例如,事务api需要Kafka版本至少0.11.0。当使用的代码和实际版本不符合时将收到unsupportedVersionException异常。

要使用事务Producer API,必须设置 transactional.id 属性, 设置transactional.id可用来在单个Producer的实例跨多个会话时方便进行事务的恢复(回滚操作)。只要 transactional.id 设置好后,幂等性所依赖的所有相关的Producer配置就会自动生效。此外,事务中包含的主题应该做好持久性设置。特别是replication.factor 应该至少是3,主题的min.insync.replicas 应设为2。最后,为了从端到端实现事务保证,还必须将使用者的事务隔离级别配置为只读取提交(read only committed messages)。

所有新的事务api都是阻塞的,一旦失败就会抛出异常。下面的示例说明了如何使用新的api。它与上面的示例类似,只是100条消息都是属于一个事务(要么同时成功,要么同时失败)。

 Properties props = new Properties();
 props.put("bootstrap.servers", "localhost:9092");
 props.put("transactional.id", "my-transactional-id");
 Producer<String, String> producer = new KafkaProducer<>(props, new StringSerializer(), new StringSerializer());

 // 初始化事务
 producer.initTransactions();

 try {
     // 事务开始
     producer.beginTransaction();
     for (int i = 0; i < 100; i++)
         producer.send(new ProducerRecord<>("my-topic", Integer.toString(i), Integer.toString(i)));
     producer.commitTransaction();
 } catch (ProducerFencedException | OutOfOrderSequenceException | AuthorizationException e) {
     // We can't recover from these exceptions, so our only option is to close the producer and exit.
     producer.close();
 } catch (KafkaException e) {
     // For all other exceptions, just abort the transaction and try again.
     // 退出事务
     producer.abortTransaction();
 }
 // 关闭生产者
 producer.close();

正如示例中所暗示的,每个生产者只能有一个开放的事务。所有在beginTransaction()commitTransaction()调用之间发送的消息都是单个事务的一部分。如果指定了事务id,则生产者发送的所有消息必须是事务的一部分。

事务Producer通过抛出异常来传递错误状态。因此不需要为producer.send()指定回调,也不需要调用producer.send().get()来返回将来执行的结果. 如果任何producer.send()或事务调用在事务期间遇到不可恢复的错误,就会抛出KafkaException。有关从事务发送中检测错误的详细信息,请参阅文档 send(ProducerRecord) 。通过在接收到KafkaException时调用producer.abortTransaction(),从而保证事务性。

Views: 499