函数式编程39:有界并发的本地订单导入工程
本章把金额计算、输入校验、不可变状态、效果描述、资源作用域和有界并发放到同一条可运行路径中。程序读取真实本地文件,按订单聚合商品行,查询模拟折扣报价,计算应付金额,产生已处理订单,并在整批成功后一次保存结果与命令回执。
它是进程内教学工具,不连接生产服务。输入、故障和调度关系由测试确定;文件确实创建、读取并关闭,效果确实执行,取消也等待任务结束。保存使用 Ref 中的内存快照,不把内存成功报告为数据库提交。完整实现自包含在本章目录,没有依赖其他章节的 Main 类。
从基线保留什么
导读与能力自测的金额样本是两件单价十二元五角的商品,加一件三元二角的商品,九折后应付二十五元三角八分。本章保留相同计算顺序:先汇总订单行金额,再应用折扣率,最后以 HALF_UP 舍入到两位小数。
这里的 discount 表示减免比例,0.10 对应九折,而不是直接乘以 0.10。零折扣得到原金额,全额折扣得到零。模拟报价端返回折扣率,因此它不是文件中单价的来源。商品单价来自输入,外部步骤决定折扣,这两个来源在模型中明确分开。
日期不参与本章结果,程序也不读取系统时钟或随机数。与第 00 章相比,这是输出模型的明确收窄,不声称已经复现带日期 Quote 的全部字段。金额规则和基线示例则作为累计回归保留下来。
本章实现没有复制旧 Scala 工程的源码,也没有导入共享 labs。它复用的是前面已经解释的设计:纯解析和状态转换、显式错误、Resource 生命周期、IO 组合以及有序的并发映射。只把当前需求需要的部分组合起来,没有为展示术语而加入额外的 Reader、Writer 或自由单子层。
五字段输入与订单行聚合
每行格式为 orderId,sku,quantity,unitPrice,currency,不带表头。订单标识和商品编码必须以大写字母开头,后续只允许大写字母与数字。数量为一到一百,单价整数部分一到七位,可有一到两位小数,币种只接受 CNY。
1 | |
前两行聚合成订单 A,后两行分别为 B 与 C。订单输出顺序按标识第一次出现的顺序保留,订单内行顺序按输入保留。相同订单下相同 sku 再次出现被判为 duplicate-line,不自动相加;不同订单允许出现相同 sku。
这个格式不是完整 CSV。它不支持引号包裹、字段内逗号或跨行字段,包含引号的标识会被拒绝。拒绝这些形式是当前协议的一部分,不能把 split(“,”) 描述为通用 CSV 解析器。需要支持标准 CSV 时,应更换解析边界并保留领域校验,而不是继续叠加零散字符串替换。
聚合也有明确的内存含义。程序先读取并校验整个有界批次,再开始报价,不声称逐行常量内存流式处理。一个批次最多一千行,每行最多两百个字符;这个上限让全批拒绝和整批提交可以在简单模型中实现。
空文件是合法的空批次,返回空向量并记录命令回执。它不构造一个没有商品行的订单,因此与“单个订单不能为空”并不冲突。空白行不是空文件,仍会进入格式校验并被拒绝。
字段错误在效果之前累计
parse 返回 Either[InputFailure,Vector[Order]]。字段数量或标识形式错误时记录 format;结构正确后独立检查数量、金额和币种,因此同一行可以产生多个错误。错误向量按行号和检查顺序输出,便于稳定重放。
输入 bad 和 A,X,0,-1,CNY 的错误期望恰好是 1:format、2:quantity、2:money。这个断言不仅检查失败,还检查错误内容和顺序。金额负数不会因为 BigDecimal 能表示它就被接受,输入域先由业务语法约束。
单价文本 1.001、NaN、1e2、10000000 均被拒绝。前者超出小数位;后两种数值表示不属于当前十进制文本协议;最后一种超过整数位上限。输入价格不会在解析阶段静默舍入,避免把用户输入的精度错误掩盖为正常金额。
解析时使用局部 LinkedHashMap 聚合,返回时转换为不可变 Vector 和 case class。局部可变容器没有逃逸到后续报价任务,后续任务不会在并发过程中向同一个聚合表追加商品行。纯边界关注的是外部可观察修改,而不是禁止函数内部使用任何可变构建器。
只要错误向量非空,整批返回 Left,已解析的合法部分也不会继续报价。这是“先验证全批再执行”的政策,适合当前原子保存模型。若希望部分成功,需要重新定义每行结果、订单聚合失败和重复命令的回执,不能仅把 Left 忽略掉。
金额计算与状态迁移是一个结果
Line 保存 sku、quantity 和 java.math.BigDecimal 单价。Order 保存订单号、行向量、CNY 币种及阶段。settle 接收订单和折扣,返回 Either[String,Paid];Paid 内含新的 Order 与最终 amount。
1 | |
代码中的 Decimal 是 java.math.BigDecimal 的导入别名。数值从十进制文本建立,不先经过 double。无 MathContext 的加法和乘法保留精确中间值,最后显式调整 scale;BigDecimal 的数值表示和 equals 还涉及 scale,因此输入单价统一为两位,输出也统一为两位。BigDecimal API
订单 A 的小计为 28.20,乘 0.90 得到 25.3800,最终为 25.38。订单 B 的 0.05 乘 0.90 为 0.045,HALF_UP 后为 0.05,这个例子实际触发舍入。订单 C 最终为 0.90。只有基线 25.38 还不足以检验半分边界,因此测试同时保留 B。
区分整单舍入与逐行舍入,需要同一订单的两行都产生半分。本实验将 ROUND 订单的 X、Y 两个不同 sku 各设为单价 0.05、数量一,通过 process 应用 0.10 折扣。两行先合计再打九折,结果为 0.09;逐行舍入再相加则为 0.10。回归测试同时断言返回金额和 saved 中的金额为 0.09,并检查回执保存了同一结果。
新阶段不能只在一个独立单元测试中出现。实际 tracked 路径直接返回 settle 产生的 Paid,process 收集并保存这些 Paid,最终打印包含 Processed 的订单。集成断言检查每个输出订单都为 Processed,并检查保存向量等于返回向量。
原 Validated 订单保持不变;把已处理订单再次传给 settle 返回 already-processed。负折扣与超过一的折扣被拒绝;零折扣和全额折扣分别得到 28.20 和 0.00。这些纯函数检查与实际导入路径共享同一个实现,没有另写一个仅供演示的计算函数。
case class 的公开构造器仍允许调用者绕过文本解析建立数据。settle 再检查当前计算所需的数量、非负单价、币种和非空行,但它不承担全部文本协议验证。本文的处理入口始终从 parse 得到订单;若把 settle 暴露成公共 API,应明确加强其构造边界,而不是声称类型已排除所有非法值。
文件读取的真实作用域
file 使用 Resource.make 打开 BufferedReader,释放时关闭并增加 closed 计数。读取与关闭放在 IO.blocking 中,文件的作用域只覆盖 load。报价开始前,文件已经读完并关闭,不让报价等待延长文件描述符寿命。
load 逐字符读取,在当前行超过两百字符时立即拒绝,避免先用 readLine 建立任意大的字符串后才检查长度。行结束时增加计数,超过一千行拒绝;支持 LF,也会去掉行末的一个 CR。长度上限按当前读取的字符计数,不是任意编码下的字节数。
这些限制不能保证整个操作系统层面没有缓冲,也不构成抵御所有恶意文件的安全库。它们给出了应用层批次和行构建器的明确边界。对于当前本地合成输入,已经足以让空间政策可观察且可测试。
成功、解析失败和容量失败都检查关闭计数。计数更新在 reader.close 成功后发生,所以它记录的是本实验中已返回的关闭调用,而不是仅仅“尝试进入 finally”。实验没有注入 close 本身失败,也没有模拟磁盘损坏;这些不在通过范围内。
Resource 让释放动作跟随使用作用域,而不是依赖调用方记住成功路径最后手动 close。Cats Effect 文档解释了 acquire、use 和 release 的作用域组合;本章把该机制落实到真实文件与活动任务计数中。Resource 文档
有序输出与并发上限分别验证
报价阶段使用 Stream.emits(orders).covary[IO].parEvalMap(2),最后 compile.toVector。选择有序版本,是因为输出应维持订单首次出现的顺序;两个报价可以重叠执行,但输出仍按 A、B、C 排列。
FS2 3.11.0 的 parEvalMap 注释和实现明确区分并发执行与有序下游输出,maxConcurrent 限制并发效果数量。本章使用冻结版本源码核对这一契约,不把新版网页行为直接当作旧依赖保证。FS2 3.11.0 Stream 源码
tracked 在每个报价任务开始时通过 Resource 增加 active,更新 peak 并累计 quoted,任务结束、失败或取消时减少 active。正常测试断言 peak 恰好为二、最终 active 为零。只检查最终零不能证明并发上限;只检查峰值也不能证明任务全部释放。
测试不用 sleep 制造重叠。A 进入后完成 entered 信号,再等待 release;B 等到 entered 后释放 A。这样 A 的完成依赖 B 已经进入,实际建立两个活动任务的重叠。这个协议能发现把有界并发误改为完全串行的实现,因为串行版本无法让 B 进入,只会被总超时看门狗终止。
输出顺序与完成顺序不同。A 可以在 B 放行之后才完成,parEvalMap 仍维持输入顺序。若改为无序版本,结果集合可能相同,但向量顺序不再符合当前契约。是否允许无序应由导入工具接口决定,不由哪个操作符看起来更快决定。
有序输出可能产生队头等待。较早订单很慢时,后续已完成结果需要等待。当前批次有容量上限,没有声称完全消除此成本;若业务只需要按订单号查询最终结果,可以另行评估无序处理,但必须修改顺序断言并解释接口变化。
整批提交与命令回执
process 接收文件路径、命令键、Audit Ref 和报价函数。读入合法订单后,先查询 receipts。如果相同键已经对应相同的规范化订单向量,直接返回此前结果,不再次报价或追加 saved;相同键对应不同输入则返回 key-conflict。
首次执行时,所有报价与 settle 都成功后才进入 Ref.modify,一次更新 saved 与 receipts。回执保存输入向量和结果向量,不只保存一个已见键。这样可以区分合法重试与错误复用请求键。
集成测试在同一个真实文件路径上调用 process 两次,第二次传入一个“只要被调用就抛 AssertionError”的报价函数。返回值必须与第一次相同,quoted 总数仍为三,saved 仍只有三个订单。这个测试经过实际入口,不是仅对独立幂等函数调用两次。
随后使用相同命令键读取不同数量的文件,断言 key-conflict,保存结果不变。三个入口调用都真实打开并关闭文件,因此最终 opened=3、closed=3。重复请求仍需读取输入才能比较规范化内容;它省略的是报价与保存,不是全部 I/O。
提交处再次检查回执,防止两个并发调用都在初始查询时看到缺席而重复追加。在同一 Ref 内,结果与回执更新是一个原子修改。这个设计不承诺并发同键请求只报价一次:两个请求在提交前仍可能分别查询外部端口,只有提交去重。报价必须适合这种重试语义。
回执相等比较的是解析后的订单向量,不是原文件字节。等价的价格文本在规范化后可以相等,商品行或订单顺序不同则可能不同。这个选择保留当前顺序契约;如果希望不考虑行顺序,必须设计稳定规范化和相应测试。
回执没有跨进程持久化,也没有过期策略。当前单次教学运行只处理有界批次,不能把 Map 的长期增长问题忽略为已解决。长期服务需要定义容量、清理和旧命令重放的规则,并将幂等状态与业务保存放在一致的耐久边界中。
故障与取消不能留下部分保存
报价失败测试让模拟端口直接抛 quote-down。process 返回失败,saved 与 receipts 为空,active 为零,文件关闭计数为一。因为保存发生在 compile.toVector 之后,任何一个报价失败都会阻止本批进入提交步骤。
混合场景用两个 Deferred 固定报价事件的先后。A 的报价先等待 B 进入,再产生折扣值;其成功分支记录 A:quote-success 并放行 B。B 随后记录 B:quote-failure,抛出 mixed-quote-down。测试断言事件序列恰为这两项,process 返回指定异常,saved 和 receipts 均为空,active 为零,opened 与 closed 均为一。成功事件来自报价动作的 Outcome.Succeeded 分支,只说明成功报价先发生,不证明 A 已完成 settle 或 tracked 已返回 Paid。门控不依赖 sleep,峰值断言仍要求两个报价任务同时占用资源。
这只保证本地保存快照不部分更新,不会撤销已经发生的外部报价调用。当前报价是受控模拟,没有扣款或发消息。如果真实端口会产生不可逆副作用,就必须改变协议,不能依靠最终向量未保存来声称整个世界已回滚。
取消测试建立 started 信号,让报价进入后执行 IO.never。测试启动 process fiber,等待 started,再调用 fiber.cancel,并等待 join。随后断言 Outcome 为取消、active 为零、saved 和 receipts 均为空,文件已经关闭。
等待 join 是关键。仅发出取消请求后立刻检查状态,可能观察到释放尚未完成的中间时刻。当前测试把“请求取消”与“取消完成并释放”分开,检查后者,避免用任意延迟猜测清理时间。
这个取消场景发生在提交之前。若调用方在 Ref.modify 已提交之后才取消等待,数据不会自动撤回;调用方可能没有收到结果,但重放相同命令可取回回执。把提交后的取消解释为回滚,会制造数据与客户端认知之间的矛盾。
总超时二十秒只是看门狗,防止同步协议写错后无限挂起。它不是性能门槛,不报告导入吞吐量,也不利用超时阈值证明调度公平性。正常、失败和取消断言使用明确事件关系与最终状态。
验收覆盖及尚未覆盖的行为
实际测试覆盖正常基线与半分舍入、整单与逐行舍入的差异、负价格、数量零与越界、金额位数、指数和 NaN 文本、错误币种、重复商品行、引号格式、空批次、超行数、超行长、报价异常、成功报价之后另一报价失败、同键重放、同键冲突,以及取消后的任务释放。各场景检查错误之外是否发生了报价、保存和回执更新。
合法批次的输出包含订单号、CNY、商品行、Processed 阶段和最终金额。程序打印的是实际返回结果,而不是预先写好的“成功样例”。源码中的 PASS 行只在前面所有断言通过后输出,runner 同时记录退出码和源码散列。
没有覆盖真实 HTTP、数据库故障、分布式锁、磁盘输出原子替换、跨进程恢复或无限输入。当前工具以运行回归入口为主要使用方式,不提供完整命令行参数解析和部署配置。这些边界需要按真实产品需求增加,不应在结课时用“工程化完成”一词抹掉。
旧文纯核心与应用边界提醒两字段输入不是完整 CSV、Future.traverse 不等于有界并发,并讨论文件求值的作用域。本章保留这些限制意识,增加五字段金额输入、订单行聚合、实际并发上限和取消完成断言,没有修改旧文或旧实验。
系列累计能力可以落在具体接口上:parse 对应纯计算与显式错误,Line 与 Order 对应不可变数据,settle 对应状态与金额规则,quote 参数对应依赖边界,Resource 对应生命周期,parEvalMap 对应有界异步,Receipt 对应重复命令协议。没有必要把每个学过的抽象都塞进一个方法。
学习覆盖表
下表把能力落实到当前工程中的观察点。跨框架一列描述迁移时仍需保留的契约,不声称本章已经运行其他语言或效果框架的对应实现。
| 能力 | 当前函数或实际证据 | 跨框架解释 |
|---|---|---|
| 纯计算与显式环境 | settle 接收订单和折扣;基线 25.38、零折扣和全额折扣断言 | 普通函数即可表达;依赖作为值传入,不能在函数内部重新读取报价 |
| 不可变数据与状态 | Order.copy 产生 Processed;最终 paid 与 saved 都保留新状态 | 新快照与旧对象就地修改是不同所有权策略;选择哪种语言都要检查旧快照未变 |
| 解析与独立校验 | parse 累计 quantity、money、currency;错误路径 quoted=0 | 可用显式结果或异常适配器,但累计错误与首次失败的政策不能无声切换 |
| 依赖端口与效果描述 | quote 参数支持正常、quote-down 和永不完成三种实现 | 普通接口也能替换依赖;是否延迟启动、是否可重复执行需另外验证 |
| 资源生命周期 | file、tracked;closed 与 active 的最终断言 | 对应作用域拥有资源的责任;语言级自动关闭仍需确认异步工作没有逃出作用域 |
| 有界并发与顺序 | parEvalMap(2)、门闩协议、peak=2 和 A/B/C 顺序 | 线程池大小、任务数与结果顺序不是同一指标;迁移后要保留每项观察 |
| 取消与任务归属 | started 后 cancel,再 join;取消后零保存与零活动任务 | 取消请求与任务结束不同;其他运行时也需要等待可观察的清理完成 |
| 重放与提交 | process 的 Receipt 与 Ref.modify;重复入口不再报价,冲突键拒绝 | 幂等属于命令和存储协议;结果句柄可重用不等于业务动作只执行一次 |
表中的证据是能力已经进入当前代码路径的依据,不是读者已经掌握能力的证明。结课检查可以从任意一行反向推导:删掉对应机制后,哪个断言应该失败?例如去掉金额规范化,输入等价与回执比较可能变化;把保存移进逐订单处理,报价失败时可能留下部分结果;只请求取消而不等待,最终计数就可能仍处于中间状态。
迁移到另一语言时,先保留这些可观察结果,再选择该生态的接口。没有必要寻找同名的 Resource 或 parEvalMap 方法才开始迁移,也不能找到同名方法后跳过验证。资源作用域、错误传播、输出顺序与命令提交是业务和执行契约,具体类型只是承载方式。
运行、手算与修改
完整自包含 Main.scala内含确定输入,临时文件由 Resource 创建并删除,不提交额外夹具或日志。运行证据记录实际命令、冻结版本、退出码、源码散列和输出。执行:
1 | |
正常结果为 A=25.38、B=0.05、C=0.90,三个新状态均为 Processed,peak=2、opened=3、closed=3、quoted=3。整单舍入用例输出 rounding returned=0.09 saved=0.09;混合失败用例记录 A:quote-success,B:quote-failure,saved=0、receipts=0、active=0、opened=1、closed=1、peak=2。非法输入、直接报价失败和取消路径也由同一次入口执行。
手算题:两行单价均为 0.05、数量均为一、折扣 0.10,逐行舍入后相加与整单汇总后舍入分别得到多少?前者 0.10,后者 0.09。当前契约采用后者,process 回归测试检查返回值与保存值均为 0.09。说明该测试为什么比只检查 25.38 更能区分舍入步骤。
修改题:把报价端改为对指定订单返回非法折扣,断言没有任何本批保存或回执。保留现有取消门闩和峰值断言,不能为了新增功能改成顺序执行或删除失败测试。
进阶修改是将 Ref 保存替换为本地文件输出。先定义提交点、临时文件清理和重复请求恢复,再选择资源与原子替换策略;验证输出确实可重新读取后才能声称保存成功。当前实验没有完成这项修改,练习要求不能作为已交付功能引用。

