Java 函数式编程常用套路(二)

从数据库捞全量数据

ScrollHelper:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
@Slf4j
public class ScrollHelper {

/**
* @param startId 从哪里开始捞取
* @param limit 每次捞取多少
* @param fetcher 怎么捞下一批(将 lastId 传入你的 Dao 层方法)
* @param idExtractor 怎么获取主键 ID(通过方法引用指定主键)
* @param consumer 捞出来这一批(1000条)后,具体做什么业务
* @param <ID>
* @param <T>
*/
public static <ID, T> void scroll(ID startId,
int limit,
BiFunction<ID, Integer, List<T>> fetcher,
Function<T, ID> idExtractor,
Consumer<List<T>> consumer) {
ID currentId = startId;
List<T> batchList;
while (!CollectionUtils.isEmpty(batchList = fetcher.apply(currentId, limit))) {
try {
// 消费掉这批数据
consumer.accept(batchList);
} catch (Exception e) {
log.error("scroll query error.", e);
} finally {
// 提取最后一项的 ID 作为下一次的起点
T lastOne = batchList.get(batchList.size() - 1);
currentId = idExtractor.apply(lastOne);
}
}
}
}

使用 ScrollHelper:

1
2
3
4
5
6
7
8
9
10
public List<UserPO> queryAllUser() {
List<UserPO> resultList = new LinkedList<>();
ScrollHelper.scroll(
0L,
1000,
(lastId, limit) -> userMapper.queryUserByCursor(lastId, limit),
UserPO::getId,
resultList::addAll);
return resultList;
}


构建一个部门员工层级树

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115

@Override
public DeptAndUserTreeResp getDeptAndUserTree() {
// 1. 批量高效获取缓存数据
List<UserPO> users = userDeptCache.queryAllUser();
List<DeptPO> depts = userDeptCache.queryAllDept();

// 将部门转化为 Map 方便快速索引
Map<Long, DeptPO> deptPOMap = depts.stream()
.collect(Collectors.toMap(DeptPO::getDeptId, v -> v, (v1, v2) -> v1));

// 统一节点容器 (预估容量提升性能)
Map<Long, UserDeptItem> itemMap = new HashMap<>(depts.size() * 2);
Set<Long> rootDeptIds = new LinkedHashSet<>(); // 使用 Set 去重,保持根节点有序

// 2. 线性时间(O(N))构建完美的部门双向树
for (DeptPO deptPO : depts) {
Long deptId = deptPO.getDeptId();
Long parentId = deptPO.getParentId();

// 2.1 填充或创建当前部门节点(完美融合提前创建的空壳)
UserDeptItem currentItem = itemMap.compute(deptId, (k, v) -> fillDeptItem(v, deptPO));

// 2.2 划分层级:判定是否为顶级根部门
if (parentId == null || !deptPOMap.containsKey(parentId)) {
rootDeptIds.add(deptId);
} else {
// 2.3 动态挂载到父节点(若父节点未创建,则初始化空壳)
itemMap.computeIfAbsent(parentId, k -> new UserDeptItem().setItems(new ArrayList<>()))
.getItems().add(currentItem);
}
}

// 3. 构建离职员工虚拟大本营
UserDeptItem leaveGroup = new UserDeptItem()
.setId(-1L)
.setName("离职")
.setType(RelationTypeEnum.DEPARTMENT.getType()) // 虚拟部门类型
.setItems(new ArrayList<>());

// 4. 将员工归类并流式挂载
for (UserPO userPO : Users) {
if (userPO.getState() == null) continue;

// 统一构建员工 DTO 项
UserDeptItem staffItem = new UserDeptItem()
.setId(userPO.getUserId())
.setName(userPO.getName())
.setEmail(userPO.getEmail())
.setNumber(userPO.getNumber())
.setType(RelationTypeEnum.STAFF.getType());

if (userPO.getState() == UserStateEnum.LEAVE.getCode()) {
leaveGroup.getItems().add(staffItem);
} else {
Long deptId = userPO.getDeptId();
UserDeptItem deptItem = itemMap.get(deptId);
if (deptItem != null) {
deptItem.getItems().add(staffItem);
} else {
log.warn("发现孤儿员工,其归属部门不存在: userId={}, deptId={}", userPO.getUserId(), deptId);
}
}
}

// 5. 聚合根节点并进行剪枝(剔除空部门)
List<UserDeptItem> finalTree = rootDeptIds.stream()
.map(itemMap::get)
.filter(Objects::nonNull)
.collect(Collectors.toList());
finalTree.add(leaveGroup);

// 深度递归移除空部门
removeEmptyDept(finalTree);

return new DeptAndUserTreeResp().setItems(finalTree);
}

/**
* 辅助方法:无缝充实或新创部门节点,保留已挂载的 items
*/
private UserDeptItem fillDeptItem(UserDeptItem existItem, DeptPO po) {
UserDeptItem item = (existItem != null) ? existItem : new UserDeptItem();
if (item.getItems() == null) {
item.setItems(new ArrayList<>());
}
item.setId(po.getDeptId());
item.setName(po.getName());
item.setLeaderId(po.getLeaderId());
item.setType(RelationTypeEnum.DEPARTMENT.getType());
return item;
}

/**
* 移除空部门
*/
private void removeEmptyDept(List<UserDeptItem> items) {
if (CollectionUtils.isEmpty(items)) {
return;
}
Iterator<UserDeptItem> iterator = items.iterator();
while (iterator.hasNext()) {
UserDeptItem next = iterator.next();
// 如果是员工节点,属于不可裁剪的“叶子实心点”,直接跳过
if (Objects.equals(next.getType(), RelationTypeEnum.STAFF.getType())) {
continue;
}
// 无论如何,先向下递归掏空它的所有子孙部门
removeEmptyDept(next.getItems());
// 回溯判定:子孙部门被掏空后,如果自己彻底变成了一个没有任何员工和子部门的裸壳,则无情剪掉
if (CollectionUtils.isEmpty(next.getItems())) {
iterator.remove();
}
}
}


数据库批量增改删数据

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
@Transactional(rollbackFor = Exception.class)
public void batchSyncEmpAssets(String employeeId, List<DeviceAssetPO> incomingList) {
if (StringUtils.isBlank(employeeId)) {
return;
}
// 如果上游传来空集合,代表该员工名下已经没有任何资产,直接全删本地资产后平账返回
if (CollectionUtils.isEmpty(incomingList)) {
List<DeviceAssetPO> oldAssets = deviceAssetDao.selectByEmployeeId(employeeId);
if (CollectionUtils.isNotEmpty(oldAssets)) {
List<Long> allIds = oldAssets.stream().map(DeviceAssetPO::getId).collect(Collectors.toList());
deviceAssetDao.batchDelete(allIds);
}
return;
}

// 1. 摘取本地旧账
List<DeviceAssetPO> persistedList = deviceAssetDao.selectByEmployeeId(employeeId);

// 2. 将旧账转化为内存 Map,以 deviceUuid 作为唯一索引联合 Key
Map<String, DeviceAssetPO> persistedMap = Optional.ofNullable(persistedList)
.orElse(Collections.emptyList())
.stream()
.collect(Collectors.toMap(DeviceAssetPO::getDeviceUuid, Function.identity(), (v1, v2) -> v2));

// 3. 准备三个分流漏斗
List<DeviceAssetPO> toInsert = new ArrayList<>();
List<DeviceAssetPO> toUpdate = new ArrayList<>();

// 4. 双向核对
for (DeviceAssetPO newItem : incomingList) {
// 顺手规范工号绑定
newItem.setEmployeeId(employeeId);

// 从旧账 Map 里剔除当前敲门的 Key
DeviceAssetPO oldItem = persistedMap.remove(newItem.getDeviceUuid());

if (oldItem == null) {
// Case A: 旧账里没有这个 UUID -> 归入「待新增」
toInsert.add(newItem);
} else {
// Case B: 旧账里有,检查非核心字段(设备名、状态)是否发生变动
boolean isChanged = !Objects.equals(newItem.getDeviceName(), oldItem.getDeviceName())
|| !Objects.equals(newItem.getStatus(), oldItem.getStatus());
if (isChanged) {
newItem.setId(oldItem.getId()); ////
toUpdate.add(newItem);
}
}
}

// 5. 收尾清算:persistedMap 中剩下的,说明上游最新的全量列表里已经没它们了 -> 属于被收回的「待删除」资产
List<Long> toDeleteIds = persistedMap.values().stream()
.map(DeviceAssetPO::getId)
.collect(Collectors.toList());

// 6. 精准组合落库,拒绝大范围无效 I/O 交互
if (CollectionUtils.isNotEmpty(toInsert)) {
int inserted = deviceAssetDao.batchInsert(toInsert);
log.info("员工 [{}] 资产同步平账 - 批量新增成功: {} 条", employeeId, inserted);
}
if (CollectionUtils.isNotEmpty(toUpdate)) {
int updated = deviceAssetDao.batchUpdate(toUpdate);
log.info("员工 [{}] 资产同步平账 - 批量修改成功: {} 条", employeeId, updated);
}
if (CollectionUtils.isNotEmpty(toDeleteIds)) {
int deleted = deviceAssetDao.batchDelete(toDeleteIds);
log.info("员工 [{}] 资产同步平账 - 批量物理清除成功: {} 条", employeeId, deleted);
}
}