Ohhnews

分类导航

$ cd ..
DZone Java原文

Jakarta Batch实战:面向企业工作负载的可靠分块处理

#jakarta batch#批处理#企业应用#分块处理#数据可靠性

批处理仍然至关重要,因为许多业务操作并不适合交互式请求。诸如重新计算价格、对账交易、迁移记录、生成报表、处理发票、重新分类客户,或对数百万条记录应用规则等任务,可能需要相当长的时间。将这些任务作为标准请求处理会导致系统脆弱、用户等待时间增加、频繁超时、重试困难以及可能的数据不一致。批处理模型能够以可预测、增量的方式处理大规模工作负载,并对进度和恢复进行控制。与其将庞大的操作作为单个循环来处理,批处理使用作业、步骤、块、检查点、过滤和可重启性。这种方法将长时间运行的数据任务与用户体验分离,同时提供结构化的执行模型。在本文中,我们将聚焦 Jakarta Batch,并考察其在现代企业应用中的持续相关性。

为什么批处理在企业系统中仍然重要

现代应用提供多种后台处理方式,例如消息队列、事件驱动架构、调度器、响应式管道和分布式流处理平台。尽管每种方式都针对特定需求,但当操作具有明确的开始和结束、涉及已知或可发现的数据集,并且需要受控执行、进度跟踪、可重启性或周期性处理时,批处理最为有效。

批处理在企业系统中仍然不可或缺。财务对账、计费、工资单、报表、数据迁移、监管处理、目录更新和大规模重新分类等工作负载仍然普遍存在。在这些场景中,首要任务是安全且可预测地处理大量工作,而不是快速响应单个事件。批处理提供了一种专门为这些需求设计的模型。

Jakarta Batch 如何工作

Jakarta Batch 将后台处理组织为作业和步骤。一个作业定义整体批处理操作,而每个步骤代表一个特定阶段。在面向块的处理中,一个步骤遵循简单的流水线:读取、处理、写入,并重复直到处理完所有输入。Jakarta Batch 运行时管理这一生命周期,因此应用代码可以专注于读取、转换和持久化数据。

[LOADING...]

作业是顶层执行单元,代表一个完整的业务操作,例如导入记录、重新计算客户分类、处理发票或对账交易。作业可以在启动时接受参数,从而允许同一批处理定义使用不同的输入或业务规则运行。一个作业由一个或多个步骤组成,每个步骤代表工作负载的一个不同阶段。简单作业可能只有一个步骤,而复杂流程可以按顺序使用多个步骤,例如导入数据、验证数据并生成最终报告。

  • 在面向块的步骤中,ItemReader 从数据库或文件等来源一次向运行时提供一个数据项。读取器只负责获取下一个数据项,不需要知道它将被如何处理或持久化。
  • ItemProcessor 接收每个数据项并应用业务规则,例如验证、转换、分类、丰富或过滤。它可以返回修改后的数据项,或者在数据项应被排除在写入之外时返回 null。
  • ItemWriter 接收已处理的数据项并持久化或导出它们。与逐个处理数据项的读取器和处理器不同,写入器通常接收来自当前块的一组数据项。这可以实现更高效的数据库或批量操作。

Jakarta Batch 在此流水线周围添加了功能,以支持企业工作负载。运行时管理块边界、事务、检查点、执行状态、失败和重启行为。块大小决定在写入和检查点之前将多少工作分组,因此它是平衡吞吐量、内存使用、数据库成本和恢复的关键参数。核心模型很直接:

作业 → 步骤 → 读取 → 处理 → 写入 → 重复

Jakarta Batch 使业务流水线保持简单,同时由运行时管理可靠、长时间运行的数据处理所需的执行机制。

示例:使用 Jakarta Batch 进行客户细分

本示例使用电子商务客户细分场景来演示 Jakarta Batch 模型。客户根据可配置的消费阈值被分配到 Bronze、Silver、Gold 和 Platinum 等层级。当阈值变化时,应用会重新评估客户群,并仅更新分类发生变化的客户。完整应用包括 MongoDB 集成、Jakarta Faces UI、预览工作流、验证和支持服务。完整源代码可在 https://github.com/soujava/mongodb-jakarta-batch 获取。本节重点关注直接参与 Jakarta Batch 执行的类。

启动批处理作业

应用通过 CustomerSegmentationService 启动批处理过程。与读取器、处理器和写入器不同,这个类不是批处理工件。相反,它是一个应用服务,从 BatchRuntime 获取 Jakarta Batch 的 JobOperator 来启动和监控作业执行。

$ java
@ApplicationScoped
public class CustomerSegmentationService {
    public static final String JOB_NAME = "customer-segmentation";
    private volatile CustomerSegmentationPolicy currentPolicy;

    // initialization and status methods omitted

    public long start(CustomerSegmentationPolicy policy) {
        if (isRunning()) {
            throw new IllegalStateException(
                "A customer segmentation batch is already running");
        }

        Properties parameters = new Properties();
        parameters.setProperty(
            CustomerSegmentationPolicy.JOB_PARAMETER, policy.toJson());

        long executionId = BatchRuntime.getJobOperator()
            .start(JOB_NAME, parameters);

        currentPolicy = policy;
        return executionId;
    }

    public boolean isRunning() {
        // implementation omitted
    }
}

这里的关键 API 是 JobOperator,Jakarta Batch 提供它作为启动、停止、重启和检查作业的接口。在此示例中,细分策略被序列化到作业参数中,以确保每次执行都获得正确的业务规则。

读取输入

第一个批处理工件 CustomerItemReader 扩展了 Jakarta Batch 的 AbstractItemReader,以实现面向块的读取器。

$ java
@Named("customerItemReader")
@Dependent
public class CustomerItemReader extends AbstractItemReader {
    @Inject
    private CustomerRepository customerRepository;

    private List customers = List.of();
    private int nextIndex;

    @Override
    public void open(Serializable checkpoint) {
        try (Stream customerStream = customerRepository.findAll()) {
            customers = customerStream
                .sorted(Comparator.comparing(Customer::getId))
                .toList();
        }

        nextIndex = checkpoint instanceof Integer index ? index : 0;
    }

    @Override
    public Customer readItem() {
        if (nextIndex >= customers.size()) {
            return null;
        }

        return customers.get(nextIndex++);
    }

    @Override
    public Serializable checkpointInfo() {
        return nextIndex;
    }
}

这些方法是 Jakarta Batch 读取器生命周期的一部分,由 AbstractItemReader 定义。open() 准备读取器,并在可用时接受之前的检查点。readItem() 向运行时提供下一个数据项;返回 null 表示没有更多输入。checkpointInfo() 报告读取器当前位置,用于检查点。为简单起见,本示例将客户加载到内存中。对于更大的工作负载,实现可以使用分页或 MongoDB 游标,而无需改变 Jakarta Batch 模型。

处理每位客户

下一个工件实现了 Jakarta Batch 的 ItemProcessor 接口。

$ java
@Named("customerTierProcessor")
@Dependent
public class CustomerTierProcessor implements ItemProcessor {
    @Inject
    @BatchProperty(
        name = CustomerSegmentationPolicy.JOB_PARAMETER)
    private String thresholdsJson;

    private CustomerSegmentationPolicy policy;

    @PostConstruct
    void initialize() {
        policy = CustomerSegmentationPolicy.fromJson(
            thresholdsJson);
    }

    @Override
    public Customer processItem(Object item) {
        if (!(item instanceof Customer customer)) {
            throw new IllegalArgumentException(
                "Expected a Customer item");
        }

        CustomerTier calculatedTier = policy.tierFor(customer.getTotalSpent());

        if (calculatedTier == customer.getTier()) {
            return null;
        }

        return Customer.builder()
            .id(customer.getId())
            .name(customer.getName())
            .totalSpent(customer.getTotalSpent())
            .tier(calculatedTier)
            .build();
    }
}

这里的 Jakarta Batch 契约很明确:ItemProcessor 定义了 processItem()。运行时会对读取器产生的每个数据项调用该方法。处理器应用细分规则,并返回转换后的客户或 null。在 Jakarta Batch 中,返回 null 有特定含义:该数据项被过滤,不会继续传给写入器。@BatchProperty 也是 Batch 集成的一部分。它接收为此作业执行定义的 thresholds 属性,允许处理器在处理开始前重建 CustomerSegmentationPolicy。

写入结果

最后一个工件扩展了 AbstractItemWriter,这是 Jakarta Batch 用于写入块的基础实现。

$ java
@Named("customerItemWriter")
@Dependent
public class CustomerItemWriter extends AbstractItemWriter {
    @Inject
    private CustomerRepository customerRepository;

    @Override
    public void writeItems(List items) {
        List customers = items.stream()
            .map(this::toCustomer)
            .toList();

        customerRepository.saveAll(customers);
    }

    private Customer toCustomer(Object item) {
        if (item instanceof Customer customer) {
            return customer;
        }

        throw new IllegalArgumentException(
            "Expected a Customer item");
    }
}

writeItems() 由继承自 AbstractItemWriter 的 Jakarta Batch 写入器契约定义。与一次接收一个数据项的处理器不同,写入器接收一个已处理数据项的集合。在这种情况下,该集合仅包含分类发生变化的客户,因为处理器已经过滤掉了其他客户。此时,流水线的 Java 组件如下:

CustomerItemReader extends AbstractItemReader
        ↓
CustomerTierProcessor implements ItemProcessor
        ↓
CustomerItemWriter extends AbstractItemWriter

这些类型将应用代码连接到 Jakarta Batch 运行时。

使用 JSL 连接各组件

Java 类定义了行为,但 Jakarta Batch 要求将读取器、处理器和写入器显式映射到每个作业。这种编排在 JSL 中描述:

ref 值直接对应于 Java 类中使用 @Named 声明的名称:

$ java
@Named("customerItemReader")
@Named("customerTierProcessor")
@Named("customerItemWriter")

因此,XML 告诉 Jakarta Batch 运行时:对于这个步骤,使用这个读取器,然后使用这个处理器,最后使用这个写入器。它还将 thresholds 作业参数映射到处理器属性中。item-count="20" 设置此示例的块大小。Jakarta Batch 协调读取和处理,根据块生命周期定期调用写入器,并建立事务和检查点边界。值 20 仅用于演示;真实应用应根据处理成本、数据库行为、事务大小、吞吐量和恢复需求来调整块大小。本文推荐这种结构:先展示类声明,然后描述继承自 Jakarta Batch 或由 Jakarta Batch 要求的生命周期方法。这种方法有助于示例讲解 API,而不是仅仅呈现孤立的方法。

结论

Jakarta Batch 的价值在于,它将大规模数据处理转变为结构化执行模型,消除了对自定义循环和临时后台逻辑的需求。通过分离读取、处理和写入,并引入作业、步骤、检查点、可重启性和面向块执行等运行时概念,它为企业应用提供了一种可预测的方法来处理涉及数千或数百万条记录的工作负载。这使实现可以专注于业务逻辑,而 Batch 运行时管理重复的执行关注点,从而使模型更易于理解、优化,并随着工作负载增长而扩展。DZone 贡献者表达的观点仅代表他们自己。