Skip to content

บทที่ 13 — Batch (ประมวลผลเป็นก้อน) + Scheduling (สั่งงานตามเวลา) — Spring Batch + @Scheduled + Quartz

← บทที่ 12: Messaging | สารบัญ | บทที่ 14: GraphQL

ระบบจริงมีงาน "ทำเป็นคาบ" หรือ "ทำของเยอะ ๆ" ตลอด:

  • ส่งใบเสร็จทุก 1 เดือน
  • คำนวณ statement ของลูกค้า 10 ล้านคน
  • ย้ายข้อมูลเก่าออกจาก hot DB
  • Reconcile payment ทุกคืน
  • ตั้งเวลา re-index search engine

ใช้ thread/loop เขียนเองได้ แต่ เปราะ:

  • App restart ตอนกำลังรัน → data state งง
  • รัน 2 instance พร้อมกัน → ทำงานซ้ำ
  • ไม่มี monitoring / progress / retry
  • 1 row พัง → ทำต่อไม่ได้

→ ใช้ tool ที่ออกแบบมาเฉพาะ

บทนี้สอน 3 ตัว:

  • @Scheduled — งานเล็ก ๆ ในตัว app
  • Quartz — scheduler enterprise-grade
  • Spring Batch — งานใหญ่ที่ต้อง chunk + retry + restart

ใช้เวลา 3-4 ชั่วโมง

🚧 โซนขั้นสูง — ข้ามได้

บทนี้เป็น บทอ้างอิงระดับลึก มือใหม่ส่วน @Scheduled (งานตามเวลาแบบง่าย) อ่านได้สบาย ส่วน Quartz + Spring Batch (งานใหญ่ระดับ enterprise) ข้ามไปก่อนได้ กลับมาเมื่อต้องประมวลผลข้อมูลจำนวนมากหรือรันงานตามเวลาในระบบหลายเครื่อง

📋 ก่อนเริ่ม + ศัพท์ปูพื้น

  • ✅ ต้องผ่าน: Spring Boot บทที่ 1–2 (REST + JPA) และเข้าใจ transaction management เบื้องต้น

ภาพรวม batch แบบง่าย ๆ: งาน batch แบ่งเป็นลำดับชั้น 3 ระดับ

  • job (งานใหญ่) — เช่น "ส่งใบแจ้งหนี้ลูกค้าทั้งหมด"

  • step (ขั้นตอนย่อยในงาน) — เช่น อ่านลูกค้า/ผู้ใช้ → คำนวณยอด → ส่งอีเมล

  • chunk (แบ่งข้อมูลทีละชิ้น) — เช่น 100 ลูกค้า/ผู้ใช้ ต่อรอบ

  • 🔑 batch (แบตช์) = ประมวลผลข้อมูล "เป็นก้อนใหญ่ทีเดียว" (ไม่ใช่ทีละ request) · scheduling (สเก-จูล-ลิ่ง) = สั่งให้งานทำงานตามเวลา/คาบที่กำหนด · job / step = งาน 1 ชิ้น / ขั้นตอนย่อยในงาน · chunk (ชังก์) = แบ่งข้อมูลเป็นชิ้นเล็ก ๆ ประมวลผลทีละชิ้น · idempotent (ไอ-เดม-โพ-เทนท์) = รันงานซ้ำได้โดยผลไม่เพี้ยน (เผื่อ job ล่มกลางคันแล้วต้องรันใหม่)


Part 1: @Scheduled — ง่ายและพอใช้

1.1 เปิด

งานที่ต้องทำตามเวลา (cleanup, report, polling) ทำได้ง่ายด้วย @Scheduled — แค่เปิดด้วย @EnableScheduling ที่ main class ก่อน แล้วจึงแปะ @Scheduled กับ method ที่อยากให้รันเป็นรอบ:

java
@SpringBootApplication
@EnableScheduling
public class App { }

1.2 รูปแบบ

@Scheduled มีหลายแบบให้เลือกตามความต้องการ — fixedRate (รันทุก N ms ไม่สนว่าครั้งก่อนเสร็จไหม), fixedDelay (รอ N ms หลังครั้งก่อนเสร็จ), cron (กำหนดเวลาแบบ crontab + timezone) เลือกให้ตรงกับลักษณะงาน:

java
@Component
public class CleanupJob {

    // ทุก 5 นาที (เริ่มหลัง app start)
    @Scheduled(fixedRate = 300_000)
    public void cleanup() { /* ... */ }

    // หลัง method เสร็จ + รอ 30 sec → รันใหม่
    @Scheduled(fixedDelay = 30_000)
    public void poll() { /* ... */ }

    // 5 sec หลัง app start ค่อยรันครั้งแรก
    @Scheduled(initialDelay = 5_000, fixedRate = 60_000)
    public void heartbeat() { /* ... */ }

    // cron expression
    @Scheduled(cron = "0 0 2 * * *", zone = "Asia/Bangkok")    // ตี 2 ทุกวัน
    public void dailyReport() { /* ... */ }
}

1.3 Cron Syntax

text
second  minute  hour  day-of-month  month  day-of-week

ตัวอย่าง:

text
0 0 2 * * *           — ตี 2 ทุกวัน
0 */15 * * * *        — ทุก 15 นาที
0 0 9 * * MON-FRI     — 9 โมงเช้า จันทร์-ศุกร์
0 0 0 1 * *           — เที่ยงคืนวันที่ 1 ของเดือน
0 30 9 ? * MON#1      — 9:30 จันทร์แรกของเดือน
                         (? = ไม่ระบุ day-of-month ใช้เมื่อระบุ day-of-week แล้ว, #1 = สัปดาห์แรกของเดือน)

⚠️ Spring cron มี 6 fields, Unix cron มี 5 fields — copy จาก crontab.guru (เว็บช่วยสร้าง cron expression สำหรับ Unix/Linux) ใส่ Spring ตรง ๆ ไม่ได้!

  • Spring cron (6 fields): second minute hour day-of-month month day-of-week — เช่น "0 0 2 * * *" = ตี 2 ทุกวัน (second=0)
  • Unix cron (5 fields): minute hour day-of-month month day-of-week — เช่น "0 2 * * *" = ตี 2 ทุกวัน (ไม่มี second)

1.4 ⚠️ ปัญหาเมื่อ scale หลาย instance

text
3 pod ใน K8s (Kubernetes — ระบบจัดการ container) + @Scheduled cron 2 AM
→ ทุก pod รันงาน → ทำซ้ำ 3 ครั้ง!

แก้:

Solution A — ShedLock

ShedLock คือ library ที่ใช้ DB เป็นตัวล็อก ป้องกันไม่ให้หลาย instance รันงานพร้อมกัน (ShedLock ไม่ได้อยู่ใน Spring Boot BOM จึงต้องระบุ version เองด้านล่าง)

xml
<!-- ShedLock ไม่ได้อยู่ใน Spring Boot BOM — ต้องระบุ version เอง
     ดู version ล่าสุดที่ https://github.com/lukas-krecan/ShedLock (last checked: 2026-06-23)
     ShedLock release บ่อย — ตรวจ version ใหม่ก่อน copy ใส่ pom.xml เสมอ -->
<dependency>
    <groupId>net.javacrumbs.shedlock</groupId>
    <artifactId>shedlock-spring</artifactId>
    <version>5.14.0</version>
</dependency>
<dependency>
    <groupId>net.javacrumbs.shedlock</groupId>
    <artifactId>shedlock-provider-jdbc-template</artifactId>
    <version>5.14.0</version>
</dependency>
sql
CREATE TABLE shedlock(
    name VARCHAR(64) PRIMARY KEY,
    lock_until TIMESTAMP WITH TIME ZONE NOT NULL,    -- PostgreSQL: timestamptz
    locked_at TIMESTAMP WITH TIME ZONE NOT NULL,     -- กัน DST (Daylight Saving Time = การปรับนาฬิกาตามฤดูกาลในบางประเทศ) / multi-timezone bug
    locked_by VARCHAR(255) NOT NULL
);

💡 ใช้ TIMESTAMP WITH TIME ZONE (Postgres: timestamptz) — TIMESTAMP ธรรมดาคำนวณ lock expiration ผิดตอนข้าม DST หรือมี pod ต่าง timezone (แม้จะใช้ usingDbTime() ก็ควรเก็บแบบ timezone-aware ไว้ก่อน)

java
@Configuration
@EnableSchedulerLock(defaultLockAtMostFor = "10m")
class SchedConfig {
    @Bean
    LockProvider lockProvider(DataSource ds) {
        return JdbcTemplateLockProvider.Configuration.builder()
            .withJdbcTemplate(new JdbcTemplate(ds))
            .usingDbTime()      // ⭐ ใช้เวลาของ DB ป้องกัน clock skew ระหว่าง pod
            .build();
    }
}

@Scheduled(cron = "0 0 2 * * *")
@SchedulerLock(name = "dailyReport", lockAtMostFor = "30m", lockAtLeastFor = "1m")
public void dailyReport() { /* รันแค่ instance เดียว */ }

⚠️ ShedLock pitfall ใน multi-instance (3+ pods)

🚧 ขั้นสูง — สำหรับ multi-instance deploy (รัน 3+ pods ขึ้นไป) ถ้า scheduler ของคุณยังรันแค่ pod เดียวข้ามหัวข้อนี้ไปก่อนได้

ก่อนอ่าน รู้จักศัพท์เพิ่ม 4 ตัว:

  • lock orphan (ล็อกค้าง) = instance ตายระหว่างถือ lock → ไม่มีใครปล่อย
  • clock skew (เวลาเหลื่อม) = นาฬิกา pod ต่างกัน 1-2 นาที เพราะ NTP (Network Time Protocol — ระบบซิงก์เวลาเครื่องให้ตรงกัน) sync ไม่พร้อมกัน
  • lost lock (ล็อกหลุด) = job รันเกิน lockAtMostForpod อื่นคิดว่า lock หมดอายุ → รันซ้อน
  • SLA (ข้อตกลงระดับบริการ) = สัญญาเวลาที่งานต้องเสร็จ เช่น payment retry ต้องเสร็จใน 30 นาที
lockAtMostFor — เวลาที่ lock ค้างได้นานสุด (กัน lock orphan)
  • ตั้งสั้นไป (เช่น 1m แต่ job รัน 30m) → instance อื่น "คิดว่า lock หมดอายุ" → รันซ้อน ❌
  • ตั้งยาวไป (24h) → ถ้าตายระหว่างรัน → ต้องรอ 24h ถึงจะ retry
  • กฎ: lockAtMostFor = max(time ของ job × 2, สั้นพอจะ retry ใน SLA)
lockAtLeastFor — เวลาที่ lock ค้างอย่างน้อย (กัน duplicate run บน clock skew)
  • แม้ job จะเสร็จเร็ว ก็ค้าง lock ไว้อย่างน้อยเท่าที่ตั้ง
  • ตั้ง > clock skew ระหว่าง pod (1-2 min พอ)
Clock skew

pod คนละเวลา → @Scheduled(cron = "0 0 2 * * *") ที่ pod-1 อาจ trigger 1 นาทีก่อน pod-2 → ทั้งสองพยายาม acquire lock → usingDbTime() แก้ได้ (ใช้เวลาของ DB เป็นแหล่งความจริงเดียว)

Lost lock

job ยาวจน lockAtMostFor หมด → pod อื่น acquire ใหม่ → 2 job รันพร้อมกัน → ✋ ใช้ idempotency table (ข้างล่าง) เพื่อกัน double-effect

Solution B — ใช้ Quartz cluster (Part 2)

Solution C — Kubernetes CronJob

แยก scheduling ออกจาก app ทั้งหมด — ให้ K8s (Kubernetes — ระบบจัดการ container, pod = หน่วยรันแอป 1 สำเนา) เป็นผู้จัดการ scheduler แทน app โดยรัน batch job เป็น pod แยกที่ตื่นมาตามเวลาที่ K8s กำหนด ข้อดี: ถ้าตายระหว่างรัน K8s retry ให้, ไม่ต้องใช้ ShedLock, K8s เก็บ history ให้

yaml
apiVersion: batch/v1
kind: CronJob
metadata:
  name: daily-report
spec:
  schedule: "0 2 * * *"
  jobTemplate:
    spec:
      template:
        spec:
          containers:
          - name: app
            image: myapp:latest
            command: ["java", "-jar", "app.jar", "--job=daily-report"]
          restartPolicy: OnFailure

💡 ฝั่ง Spring Boot ต้องตั้ง spring.main.web-application-type=none และใช้ ApplicationRunner ที่อ่าน --job=daily-report แล้วเรียก JobLauncher.run(...) จากนั้น context จะปิดเองและ pod จะ exit — ถ้าไม่ทำแบบนี้ pod จะค้างอยู่โดยไม่จบงาน

→ K8s จัดการ schedule, retry, history ให้

1.5 TaskScheduler — programmatic

@Scheduled กำหนดเวลาตายตัวตอน compile — ถ้าต้องการ schedule แบบ dynamic (รู้เวลาตอน runtime เช่น user ตั้งเอง) ใช้ TaskScheduler สั่ง schedule ผ่านโค้ดได้ ยืดหยุ่นกว่าแต่ต้องจัดการ lifecycle เอง:

java
// TaskScheduler จะถูก auto-configure ให้อัตโนมัติเมื่อใช้ @EnableScheduling
// ใน production ควรใช้ constructor injection แทน @Autowired แบบ field
@Service
public class ReminderService {
    private final TaskScheduler scheduler;

    public ReminderService(TaskScheduler scheduler) {
        this.scheduler = scheduler;
    }

    public void scheduleOnce(Instant at) {
        scheduler.schedule(
            () -> System.out.println("once"),
            at
        );
    }

    public void scheduleRepeating() {
        scheduler.scheduleAtFixedRate(
            () -> System.out.println("every 5 sec"),
            Duration.ofSeconds(5)
        );
    }
}

ใช้ตอน schedule แบบ dynamic (รู้เวลาตอน runtime)


Part 2: Quartz — Scheduler Enterprise

2.1 เมื่อไหร่ใช้ Quartz แทน @Scheduled

  • Schedule persistent (เก็บ DB) — restart app ก็ไม่หาย
  • Cluster mode — หลาย instance + lock อัตโนมัติ
  • ต้องการ trigger ที่ซับซ้อน (cron + calendar exclusion + chain)
  • Job ที่ schedule dynamic (เช่น user ตั้งเอง)
  • Misfire policy (ถ้า server down ตอนถึงเวลา — ทำยังไง)

2.2 Setup

ตั้ง Quartz ด้วย starter-quartz — จุดสำคัญที่ต่างจาก @Scheduled คือ job-store-type: jdbc (เก็บ schedule ใน DB ไม่หายเมื่อ restart) และ isClustered: true (หลาย instance ไม่รันงานซ้ำ เพราะ lock ผ่าน DB) นี่คือเหตุผลหลักที่ใช้ Quartz:

xml
<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-quartz</artifactId>
</dependency>
yaml
spring:
  quartz:
    job-store-type: jdbc           # หรือ memory (dev)
    jdbc:
      initialize-schema: never     # ใช้ script จาก Quartz เอง
    properties:
      org.quartz:
        scheduler:
          instanceName: MyAppScheduler
          instanceId: AUTO
        jobStore:
          isClustered: true        # ← cluster mode
          clusterCheckinInterval: 10000
          driverDelegateClass: org.quartz.impl.jdbcjobstore.PostgreSQLDelegate
        threadPool:
          threadCount: 10

⚠️ Quartz cluster mode = หนัก DB: cluster mode เขียน/อ่าน ~10 table ใน DB เดียวกับ app และใช้ SELECT FOR UPDATE หนัก — สำหรับ trigger ที่ถี่กว่า 1 นาที มัก saturate DB ใช้ K8s CronJob หรือแยก DB instance สำหรับ Quartz แทน

Run Quartz schema (ต้อง run manual เพราะ initialize-schema: never) ทำตามขั้นตอนนี้:

ขั้นตอนที่ 1 — ดาวน์โหลด SQL script:

ทางเลือก — ดึงจาก JAR ที่ Maven download มาแล้ว:

bash
# ค้นหา JAR ใน Maven local repository (เปลี่ยน version ตาม Quartz ที่ใช้)
find ~/.m2 -name "quartz-*.jar" | head -1
# แตก SQL script ออกมา
jar xf ~/.m2/repository/org/quartz-scheduler/quartz/<version>/quartz-<version>.jar \
    org/quartz/impl/jdbcjobstore/tables_postgres.sql

ขั้นตอนที่ 2 — Run script ใน DB:

bash
# PostgreSQL
psql -h localhost -U myuser -d mydb -f tables_postgres.sql

# หรือถ้าใช้ Docker
docker exec -i postgres_container psql -U myuser -d mydb < tables_postgres.sql

💡 Script จะสร้างตารางประมาณ 11 ตาราง (QRTZ_JOB_DETAILS, QRTZ_TRIGGERS, ฯลฯ) — ถ้ารัน script ซ้ำจะ error เพราะตารางมีอยู่แล้ว ใช้ DROP TABLE IF EXISTS ก่อนได้ถ้าต้องการ reset

2.3 Job

ใน Quartz งานหนึ่งหน่วยคือ Job (implement interface) ที่รับ parameter ผ่าน JobDataMap — แยกชัดระหว่าง "งานทำอะไร" (Job) กับ "ทำเมื่อไหร่" (Trigger ในหัวข้อถัดไป) ทำให้ reuse job เดียวกับหลาย schedule ได้:

java
// ⚠️ Quartz Job ต้องใช้ @Autowired field injection เท่านั้น — ใช้ constructor injection ไม่ได้
// เพราะ Quartz instantiate Job class ตรงๆ ด้วย default constructor ก่อน แล้ว SpringBeanJobFactory
// (ที่ spring-boot-starter-quartz ตั้งให้อัตโนมัติ) ค่อย inject dependency ผ่าน field ทีหลัง
public class SendInvoiceJob implements Job {
    @Autowired InvoiceService invoiceService;

    @Override
    public void execute(JobExecutionContext context) {
        JobDataMap data = context.getJobDetail().getJobDataMap();
        Long customerId = data.getLong("customerId");
        invoiceService.send(customerId);
    }
}

2.4 Trigger + Schedule

java
@Configuration
public class QuartzConfig {

    @Bean
    public JobDetail invoiceJobDetail() {
        return JobBuilder.newJob(SendInvoiceJob.class)
            .withIdentity("invoiceJob", "billing")
            .usingJobData("customerId", 123L)   // ใส่ data ที่ Job จะอ่านผ่าน JobDataMap
            .storeDurably()
            .build();
    }

    @Bean
    public Trigger invoiceTrigger(JobDetail invoiceJobDetail) {
        return TriggerBuilder.newTrigger()
            .forJob(invoiceJobDetail)
            .withIdentity("invoiceTrigger")
            .withSchedule(CronScheduleBuilder.cronSchedule("0 0 1 1 * ?")
                .inTimeZone(TimeZone.getTimeZone("Asia/Bangkok"))
                .withMisfireHandlingInstructionFireAndProceed())
            .build();
    }
}

Misfire policy:

  • withMisfireHandlingInstructionDoNothing — ข้าม
  • withMisfireHandlingInstructionFireAndProceed — ยิง 1 ครั้งทันทีแล้วทำตามเวลาปกติ
  • withMisfireHandlingInstructionIgnoreMisfires — ยิงทุก missed

2.5 Dynamic Schedule

จุดเด่นของ Quartz คือ schedule แบบ dynamic ที่ persist ได้ — สร้าง JobDetail + Trigger ตอน runtime (เช่น user ตั้ง reminder เอง) แล้ว schedule เก็บลง DB อยู่รอด restart ทำสิ่งที่ @Scheduled (เวลาตายตัว) ทำไม่ได้:

java
@Service
public class SchedulingService {
    private final Scheduler scheduler;

    public SchedulingService(Scheduler scheduler) {
        this.scheduler = scheduler;
    }

    public void scheduleReminder(Long userId, Instant at) throws SchedulerException {
        JobDetail job = JobBuilder.newJob(ReminderJob.class)
            .withIdentity("reminder-" + userId, "reminders")
            .usingJobData("userId", userId)
            .build();

        Trigger trigger = TriggerBuilder.newTrigger()
            .startAt(Date.from(at))
            .build();

        // กัน duplicate: ถ้า job นี้มีอยู่แล้ว ให้ reschedule แทน schedule ใหม่
        if (scheduler.checkExists(job.getKey())) {
            scheduler.rescheduleJob(trigger.getKey(), trigger);
        } else {
            scheduler.scheduleJob(job, trigger);
        }
    }
}

2.6 Listener

อยากรู้ว่า job เริ่ม/จบ/พังเมื่อไหร่ (เพื่อ log, alert, metric) ใช้ JobListener ที่มี hook ก่อน/หลัง execute — แยก cross-cutting concern (logging, monitoring) ออกจาก business logic ของ job:

java
public class JobLogListener implements JobListener {
    public String getName() { return "log-listener"; }

    public void jobToBeExecuted(JobExecutionContext ctx) {
        log.info("starting: {}", ctx.getJobDetail().getKey());
    }
    public void jobWasExecuted(JobExecutionContext ctx, JobExecutionException ex) {
        if (ex != null) log.error("failed", ex);
        else log.info("done: {} ({} ms)",
            ctx.getJobDetail().getKey(), ctx.getJobRunTime());
    }
    public void jobExecutionVetoed(JobExecutionContext ctx) { }
}

Part 3: Spring Batch — งานใหญ่

3.1 ปัญหาที่ Spring Batch แก้

ต้องประมวลผล 10M row:

  • Read ครั้งเดียวเข้า memory → OOM (Out Of Memory — หน่วยความจำเต็ม)
  • ใช้ pagination + loop เอง → ไม่มี retry / restart / skip
  • ไม่มี progress + audit trail
  • ผิดที่ row 5M → start ใหม่ตั้งแต่ 0!

Spring Batch ให้:

  • Chunk-oriented processing — read N rows → process → write
  • Restart — ผิดที่ row 5M → restart จาก 5M ได้
  • Skip + retry — row พังก็ข้าม
  • Job metadata (job_execution table) — track ทุก run
  • Parallelism — partition + multi-thread step

3.2 Concept

Spring Batch มีโครงสร้างเป็นลำดับชั้น — Job ประกอบด้วยหลาย Step, แต่ละ Step ทำงานแบบ chunk: ItemReader (อ่านทีละก้อน) → ItemProcessor (แปลง) → ItemWriter (เขียน) เข้าใจ 3 ส่วนนี้คือเข้าใจ Spring Batch ทั้งหมด:

3.3 Setup

เพิ่ม starter-batch — จุดที่ต้องรู้คือ Spring Batch ต้องมี DB เก็บ job metadata (ประวัติการรัน, สถานะ, ทำให้ restart ได้) และมัก disable job.enabled เพื่อไม่ให้ job รันเองตอน start (ให้ trigger เองเมื่อต้องการ):

xml
<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-batch</artifactId>
</dependency>

Spring Batch ต้องมี DB เก็บ job metadata:

yaml
spring:
  batch:
    jdbc:
      initialize-schema: always   # dev เท่านั้น — ⚠️ always จะ recreate schema ทุก restart → ลบ job execution history หมด
                                  # production ใช้ never (run script เอง) เพื่อให้ restart-from-checkpoint ทำงานได้
                                  # ⚠️ ถ้าต้องการทดสอบ restart-from-checkpoint ในเครื่อง dev ให้เปลี่ยนเป็น never
                                  #    แล้ว run Batch schema script เองครั้งเดียว (ดูวิธีดาวน์โหลดใน §2.2)
                                  #    มิฉะนั้น ทุก restart จะลบ execution history → test restart behavior ไม่ได้
    job:
      enabled: false              # ไม่ run job อัตโนมัติ

💡 ถ้า DB เดียวกันมีหลาย tenant/แอปใช้ schema ร่วมกัน ระวัง always เพิ่มเติม — Spring Batch 5 ใช้ property นี้ตรวจด้วยว่ามี schema เดิมอยู่แล้วหรือไม่ การรัน always ซ้ำในสภาพแวดล้อมที่ schema ถูกแชร์กับตารางแอปอื่นอาจชน permission หรือกระทบตารางอื่นได้ ไม่ใช่แค่ลบ job history เฉย ๆ — แนะนำแยก schema/datasource เฉพาะสำหรับ Spring Batch metadata

3.4 Example — ETL User → Statement

ตัวอย่างจริง — ETL (Extract-Transform-Load = ดึงข้อมูล → แปลง → โหลด) ที่อ่าน user จาก DB, สร้าง statement, แล้วเขียนกลับ ประกอบ Reader (JdbcCursorItemReader อ่านทีละ row ไม่โหลดเข้า memory ทั้งหมด) + Processor + Writer เข้าเป็น Step เห็นว่า chunk processing จัดการข้อมูลใหญ่ได้โดยไม่ OOM:

java
public record User(long id, String name, BigDecimal balance) { }
public record Statement(long userId, String pdf, BigDecimal balance) { }

Reader

java
@Bean
public JdbcCursorItemReader<User> userReader(DataSource ds) {
    return new JdbcCursorItemReaderBuilder<User>()
        .name("userReader")
        .dataSource(ds)
        .sql("SELECT id, name, balance FROM users WHERE active = true")
        .rowMapper((rs, rn) -> new User(
            rs.getLong("id"),
            rs.getString("name"),
            rs.getBigDecimal("balance")))
        .build();
}

Processor

java
@Bean
public ItemProcessor<User, Statement> processor() {
    return user -> {
        String pdf = generatePdf(user);          // method สมมติ — แทนด้วย logic สร้าง PDF ของคุณ
        return new Statement(user.id(), pdf, user.balance());
    };
}

Writer

java
@Bean
public JdbcBatchItemWriter<Statement> writer(DataSource ds) {
    return new JdbcBatchItemWriterBuilder<Statement>()
        .dataSource(ds)
        .sql("INSERT INTO statements (user_id, pdf, balance) VALUES (:userId, :pdf, :balance)")
        .beanMapped()   // ใช้ BeanPropertySqlParameterSource — map Java record accessor (userId(), pdf(), balance()) ให้อัตโนมัติ
        .build();
}

Step

java
@Bean
public Step generateStatementsStep(
        JobRepository jobRepo,
        PlatformTransactionManager tx,
        JdbcCursorItemReader<User> reader,
        ItemProcessor<User, Statement> processor,
        JdbcBatchItemWriter<Statement> writer) {

    return new StepBuilder("generateStatements", jobRepo)
        .<User, Statement>chunk(100, tx)               // chunk size = 100
        .reader(reader)
        .processor(processor)
        .writer(writer)
        .faultTolerant()
        .skipLimit(10)
        .skip(IllegalArgumentException.class)
        .retryLimit(3)
        .retry(org.springframework.dao.TransientDataAccessException.class)  // ใช้ exception ที่มีจริงใน Spring — แทนด้วย exception ชั่วคราวของระบบคุณ
        .build();
}

💡 หลาย TransactionManager ใน project เดียว: ถ้า project มีทั้ง JPA และ JDBC datasource Spring จะพบ PlatformTransactionManager หลายตัวและ throw NoUniqueBeanDefinitionException ตอน start
แก้โดย: (1) แปะ @Primary บน transaction manager ที่ต้องการให้ batch ใช้ หรือ (2) inject ด้วยชื่อเฉพาะ @Qualifier("batchTransactionManager") แล้วส่งเข้า StepBuilder ตรงๆ

Job

java
@Bean
public Job statementJob(JobRepository jobRepo, Step generateStatementsStep) {
    return new JobBuilder("statementJob", jobRepo)
        .start(generateStatementsStep)
        .incrementer(new RunIdIncrementer())  // ⚠️ ตั้งใจให้รันใหม่ทุกครั้ง — ดูข้อควรระวังเรื่อง restart ใน 3.5
        .build();
}

Run

java
@RestController
class JobController {
    private final JobLauncher launcher;
    private final Job statementJob;

    public JobController(JobLauncher launcher, Job statementJob) {
        this.launcher = launcher;
        this.statementJob = statementJob;
    }

    @PostMapping("/jobs/statement")
    public String run() throws Exception {
        JobParameters params = new JobParametersBuilder()
            .addLong("time", System.currentTimeMillis())
            .toJobParameters();
        var exec = launcher.run(statementJob, params);
        return exec.getStatus().toString();
    }
}

3.5 Restart

🟡 Concept สำคัญ — JobInstance vs JobExecution:

  • JobInstance = identity ของ "ครั้งที่จะรัน" = (jobName + JobParameters hash) — hash คือค่าที่คำนวณจาก parameter ทั้งหมดรวมกัน ถ้า parameter เปลี่ยนแม้แต่ตัวเดียว hash ก็เปลี่ยนตาม ทำให้ Spring Batch มองว่าเป็นคนละ instance กันทันที
  • JobExecution = ครั้งที่รันจริง 1 instance อาจมีหลาย execution ตัวอย่าง: รันครั้งแรก fail แล้ว restart ด้วย parameter เดิม → ครั้งที่สองนี้นับเป็น execution ใหม่ แต่ยังอยู่ภายใต้ instance เดิม
  • checkpoint = จุดบันทึกสถานะที่ Spring Batch เก็บไว้ใน DB บอกว่าประมวลผลไปถึงแถวที่เท่าไหร่แล้ว — ถ้า job ล้มกลางทาง ครั้งถัดไปจะ restart ต่อจากจุดนั้นได้เลยโดยไม่ต้องเริ่มใหม่จาก 0
  • Restart-from-checkpoint ทำงานได้ก็เมื่อ JobParameters เหมือนเดิม เพราะ Spring Batch ค้น BATCH_JOB_EXECUTION_CONTEXT ของ instance เดียวกัน

🔴 RunIdIncrementer ขัดกับ restart: ใส่ .incrementer(new RunIdIncrementer()) ทุกรอบจะเปลี่ยน run.id parameter → JobInstance ใหม่ ทุกครั้ง → restart-from-checkpoint ไม่ทำงาน เพราะ instance ใหม่เริ่มจาก 0 เสมอ

เลือกใช้ตามเจตนา:

  • อยาก restart-from-checkpoint ได้ → อย่าใส่ incrementer; รัน fail แล้วเรียกซ้ำด้วย JobParameters เดิม
  • อยาก force ให้รันใหม่ทุกครั้ง (idempotent job) → ใส่ RunIdIncrementer หรือ pass timestamp parameter
  • กรณีจริงในงาน: daily report = RunIdIncrementer (รันใหม่ทุกวัน), ETL ที่ fail ตอน restart ต้องต่อ = ไม่ใส่ incrementer

ถ้า job fail — รัน job เดิม + parameter เดิม (และไม่ใช้ incrementer) → Spring Batch จะ continue จากจุดที่หยุด

อยาก force restart from scratch (สร้าง JobInstance ใหม่ด้วย parameter ที่ไม่ซ้ำ — execution เก่าปล่อยทิ้งไว้ใน history):

java
JobParameters params = new JobParametersBuilder()
    .addLong("run.id", System.currentTimeMillis())   // unique parameter → JobInstance ใหม่
    .toJobParameters();

3.6 Parallel Processing

Multi-thread Step

🔴 กับดัก thread-safety: JdbcCursorItemReader (และ cursor reader ส่วนใหญ่) ไม่ thread-safe (thread-safe = ปลอดภัยเมื่อหลาย thread ใช้พร้อมกัน) ใช้กับ SimpleAsyncTaskExecutor ตรง ๆ จะอ่านข้อมูลซ้ำ/ข้าม/พัง — ต้องครอบด้วย SynchronizedItemStreamReader หรือสลับไปใช้ JdbcPagingItemReader ที่ stateless (ไม่เก็บสถานะภายใน อ่านแล้วไม่มีผลกระทบต่อ thread อื่น)

java
// ❌ พัง — JdbcCursorItemReader ไม่ thread-safe
.reader(cursorReader)
.taskExecutor(new SimpleAsyncTaskExecutor())

// ✅ ทางที่ 1 — ครอบ reader เดิมด้วย SynchronizedItemStreamReader
@Bean
public SynchronizedItemStreamReader<User> safeReader(JdbcCursorItemReader<User> raw) {
    SynchronizedItemStreamReader<User> sync = new SynchronizedItemStreamReader<>();
    sync.setDelegate(raw);
    return sync;
}

// ✅ ทางที่ 2 (แนะนำ) — ใช้ JdbcPagingItemReader ที่ stateless thread-safe โดย design
@Bean
public JdbcPagingItemReader<User> pagingReader(DataSource ds) {
    return new JdbcPagingItemReaderBuilder<User>()
        .name("pagingReader")
        .dataSource(ds)
        .selectClause("SELECT id, name, balance")
        .fromClause("FROM users")
        .whereClause("WHERE active = true")
        .sortKeys(Map.of("id", Order.ASCENDING))
        .pageSize(1000)
        .rowMapper((rs, i) -> new User(rs.getLong("id"), rs.getString("name"), rs.getBigDecimal("balance")))
        .build();
}

// ThreadPoolTaskExecutor bean — ตั้ง pool size ให้ชัดเจน (กำหนด concurrency จากที่นี่)
@Bean
public ThreadPoolTaskExecutor batchTaskExecutor() {
    ThreadPoolTaskExecutor exec = new ThreadPoolTaskExecutor();
    exec.setCorePoolSize(4);
    exec.setMaxPoolSize(8);
    exec.setThreadNamePrefix("batch-");
    exec.initialize();
    return exec;
}

// Step config สำหรับ multi-thread
@Bean
public Step multiThreadStep(JobRepository jobRepo, PlatformTransactionManager tx,
                            JdbcPagingItemReader<User> pagingReader,
                            ItemProcessor<User, Statement> processor,
                            JdbcBatchItemWriter<Statement> writer,
                            ThreadPoolTaskExecutor batchTaskExecutor) {
    return new StepBuilder("multiThreadStep", jobRepo)
        .<User, Statement>chunk(100, tx)
        .reader(pagingReader)
        .processor(processor)
        .writer(writer)
        .taskExecutor(batchTaskExecutor)  // ThreadPoolTaskExecutor — อย่าใช้ SimpleAsyncTaskExecutor ใน production
        // Spring Batch 5: ลบ .throttleLimit() ออกแล้ว — ควบคุม concurrency ผ่าน maxPoolSize ของ ThreadPoolTaskExecutor แทน
        .build();
}

⚠️ การ inject ด้านบนอาศัยชื่อ parameter (batchTaskExecutor) ตรงกับชื่อ @Bean method พอดี ถ้า project มี TaskExecutor/ThreadPoolTaskExecutor bean อื่นอยู่แล้ว (เช่นจาก @EnableAsync auto-config ที่สร้าง applicationTaskExecutor ให้) แล้วมีคนเปลี่ยนชื่อ parameter หรือชื่อ bean โดยไม่ตั้งใจให้ตรงกัน Spring จะหา bean ชนิดเดียวกันได้มากกว่า 1 ตัวและ throw NoUniqueBeanDefinitionException ตอน start — ป้องกันได้ด้วยการใส่ @Qualifier("batchTaskExecutor") ที่ parameter ให้ชัดเจนแทนที่จะพึ่งชื่อ

⚠️ SimpleAsyncTaskExecutor สร้าง thread ใหม่ทุก task = ไม่มี pool, ไม่มี limit จริง — ใช้ ThreadPoolTaskExecutor ที่ตั้ง corePoolSize/maxPoolSize ชัดเจน

Partitioning

แบ่งข้อมูลเป็นหลายส่วน (slice) แล้วประมวลผลพร้อมกัน (parallel) — Spring Batch 5 (Spring Boot 3.x) ลบ StepBuilderFactory แล้ว ใช้ StepBuilder ตรง ๆ + inject JobRepository แทน:

java
@Bean
public Step partitionedStep(JobRepository jobRepo,
                            PlatformTransactionManager tx,
                            Step workerStep,
                            TaskExecutor taskExecutor) {
    return new StepBuilder("partitioned", jobRepo)
        .partitioner("worker", new ColumnRangePartitioner())  // ColumnRangePartitioner มีใน Spring Batch core ตั้งแต่ 4.x
                                                               // import: org.springframework.batch.item.database.support.ColumnRangePartitioner
                                                               // constructor: new ColumnRangePartitioner() แล้ว set column/dataSource/table ก่อนใช้ (ดูตัวอย่างด้านล่าง)
        .step(workerStep)
        .gridSize(4)
        .taskExecutor(taskExecutor)
        .build();
}

แต่ละ partition รัน worker ตัวเอง (เช่น id 1-25%, 26-50%, ฯลฯ) — partition ละ JobExecutionContext แยกกัน อย่า share mutable state (สถานะที่แก้ไขได้ เช่น field ของ object ที่เปลี่ยนค่าได้) ระหว่าง worker เพราะถ้าหลาย worker เข้าถึง/แก้พร้อมกันโดยไม่ระวังจะเกิด race condition (การชิงแก้ข้อมูลพร้อมกันจนผลลัพธ์ผิดเพี้ยน)

💡 ตัวอย่าง setup ColumnRangePartitioner: (import org.springframework.batch.item.database.support.ColumnRangePartitioner)

java
@Bean
public ColumnRangePartitioner partitioner(DataSource ds) {
    ColumnRangePartitioner p = new ColumnRangePartitioner();
    p.setColumn("id");              // column ที่ใช้แบ่ง range (ต้องเป็น numeric)
    p.setTable("users");            // ชื่อ table
    p.setDataSource(ds);
    return p;
}

Partitioner จะ query MIN(id) และ MAX(id) แล้วแบ่งเป็น gridSize ranges ให้อัตโนมัติ แต่ละ partition จะได้ minValue / maxValue ใน StepExecutionContext ให้ reader ใช้กรองข้อมูล

3.7 ItemReader Types

ไม่ต้องเขียน reader เอง — Spring Batch มี ItemReader สำเร็จสำหรับแหล่งข้อมูลหลายแบบ (DB cursor/paging, JPA, CSV, XML, Kafka) เลือกตามว่า data มาจากไหน ข้อควรรู้: paging reader ปลอดภัยกว่า cursor เมื่อข้อมูลถูกเพิ่มระหว่าง read:

  • JdbcCursorItemReader — cursor SQL, stream
  • JdbcPagingItemReader — pagination (ดีกว่าตอน DB เพิ่ม row ระหว่าง read)
  • JpaPagingItemReader — ใช้ JPA
  • FlatFileItemReader — CSV/TXT
  • StaxEventItemReader — XML
  • KafkaItemReader — Kafka topic

3.8 ItemWriter Types

เช่นเดียวกับ reader — มี ItemWriter สำเร็จสำหรับปลายทางหลายแบบ (JDBC batch, JPA, CSV, Kafka) ที่ใช้บ่อยคือ JdbcBatchItemWriter (เขียนเป็น batch จึงเร็ว) และ CompositeItemWriter (เขียนหลายปลายทางพร้อมกัน):

  • JdbcBatchItemWriter — JDBC batch (เร็ว)
  • JpaItemWriter — JPA
  • FlatFileItemWriter — CSV/TXT
  • KafkaItemWriter — publish ไป Kafka
  • CompositeItemWriter — รวมหลาย writer

3.9 Listener — hook เข้า job lifecycle

Spring Batch มี listener ทุกระดับ:

Listenerhook ที่ได้
JobExecutionListenerbeforeJob / afterJob
StepExecutionListenerbeforeStep / afterStep
ChunkListenerbeforeChunk / afterChunk / afterChunkError
ItemReadListener<T>beforeRead / afterRead / onReadError
ItemProcessListener<I,O>beforeProcess / afterProcess / onProcessError
ItemWriteListener<T>beforeWrite / afterWrite / onWriteError
SkipListener<I,O>onSkipInRead / onSkipInProcess / onSkipInWrite
RetryListeneropen / onError / close
java
public class JobMonitor implements JobExecutionListener {
    public void beforeJob(JobExecution exec) {
        log.info("start job: {}", exec.getJobInstance().getJobName());
    }
    public void afterJob(JobExecution exec) {
        log.info("done: status={} read={} write={}",
            exec.getStatus(),
            exec.getStepExecutions().iterator().next().getReadCount(),
            exec.getStepExecutions().iterator().next().getWriteCount());
    }
}

// chunk-level
public class ChunkMonitor implements ChunkListener {
    public void afterChunkError(ChunkContext ctx) {
        log.error("chunk failed at step {}", ctx.getStepContext().getStepName());
    }
}

// register (ใช้ User, Statement ตามตัวอย่าง ETL ข้างบน — แทน U, S ด้วย type จริงของโปรเจกต์คุณ)
@Bean
public Step etlStep(JobRepository repo, PlatformTransactionManager tx,
                    ItemReader<User> reader, ItemProcessor<User, Statement> proc, ItemWriter<Statement> writer) {
    return new StepBuilder("etl", repo)
        .<User, Statement>chunk(100, tx)
        .reader(reader).processor(proc).writer(writer)
        .listener(new ChunkMonitor())
        .listener(new SkipListener<User, Statement>() {
            public void onSkipInProcess(User item, Throwable e) { log.warn("skip: {}", item); }
            public void onSkipInWrite(Statement item, Throwable e) {}
            public void onSkipInRead(Throwable e) {}
        })
        .faultTolerant()
            .skipLimit(10)
            .skip(jakarta.validation.ValidationException.class)  // Spring Boot 3.x ใช้ jakarta.* ไม่ใช่ javax.*
            .retryLimit(3)
            .retry(org.springframework.dao.TransientDataAccessException.class)
            .build();
}

Skip vs Retry — ต่างกันยังไง

SkipRetry
ทำอะไรข้าม item นี้ → process ตัวถัดไปลอง item เดิมซ้ำ (n ครั้ง)
ใช้กับ exceptionvalidation error, data corrupt — แก้ไม่ได้ข้อผิดพลาดชั่วคราว (transient error) เช่น network blip, deadlock
limitskipLimit(n) — เกิน → fail jobretryLimit(n) — เกิน → ส่งให้ skip (ถ้ามี)

Skip ใน Processor vs Writer — สำคัญ

java
.faultTolerant()
.skip(jakarta.validation.ValidationException.class)       // throw จาก processor → skip OK (atomic)
.skip(jakarta.validation.ConstraintViolationException.class) // throw จาก writer → ⚠️ rollback whole chunk แล้ว re-process ทีละ item

ปัญหา: writer fail = rollback ทั้ง chunk → Spring Batch re-process ทีละ item เพื่อหาตัวที่ผิด → ช้ามาก ถ้า chunk ใหญ่
ทางแก้: ตรวจ validation ใน processor ให้หมด → ส่งเฉพาะ "ของดี" ไป writer


Part 4: เลือก Tool อันไหน

Rule of thumb:

  • Heartbeat / cleanup / cache refresh → @Scheduled + ShedLock
  • User-scheduled reminder / dynamic schedule → Quartz
  • ETL / report generation / data migration → Spring Batch
  • เรียกจาก scheduler ภายนอก (Airflow, Argo Workflow) → Spring Boot CLI + K8s CronJob

Part 5: Production Best Practices

5.1 Idempotent Job

Job ต้องรันซ้ำได้โดย result เดียวกัน:

  • ใช้ unique key ก่อน insert
  • ใช้ "processed_at" column ใน source
  • Use upsert (INSERT ... ON CONFLICT UPDATE)

Deduplication Table — pattern มาตรฐาน

ใช้ตอนที่ "ทำงานชิ้นนี้ซ้ำสองครั้ง = ผิด" เช่น ส่งเงิน, ส่ง email, charge credit card

sql
CREATE TABLE processed_jobs (
    idempotency_key VARCHAR(128) PRIMARY KEY,
    job_name VARCHAR(64) NOT NULL,
    processed_at TIMESTAMPTZ DEFAULT now(),
    result JSONB
);

ใช้ใน job:

java
@Component
@RequiredArgsConstructor  // (@RequiredArgsConstructor — Lombok: สร้าง constructor จาก field final ให้อัตโนมัติ)
public class PaymentRetryJob {
    private final JdbcTemplate jdbc;
    private final PaymentApi payment;

    @Scheduled(cron = "0 */5 * * * *")
    @SchedulerLock(name = "paymentRetry", lockAtMostFor = "10m")
    @Transactional   // ต้องมี DataSource + PlatformTransactionManager ใน context (Spring Boot auto-configure ให้ถ้ามี spring-boot-starter-data-jpa)
                     // ถ้าไม่มี จะ error ตอน runtime: "No qualifying bean of type 'PlatformTransactionManager'"
    public void retryFailedPayments() {
        List<Payment> pending = jdbc.query(
            "SELECT * FROM payments WHERE status = 'FAILED' AND retry_count < 3 LIMIT 100 FOR UPDATE SKIP LOCKED",
            paymentRowMapper);

        for (Payment p : pending) {
            String key = "payment-retry:" + p.id() + ":" + p.retryCount();
            // ⚠️ ใน PostgreSQL ถ้า catch DuplicateKeyException ภายใน @Transactional เดียวกัน
            // transaction จะเข้าสู่ aborted state — query ถัดไปจะ fail ด้วย "current transaction is aborted"
            // ใช้ ON CONFLICT DO NOTHING แล้วตรวจจำนวน row ที่ insert ได้แทน
            int inserted = jdbc.update(
                "INSERT INTO processed_jobs(idempotency_key, job_name) VALUES (?, ?) ON CONFLICT DO NOTHING",
                key, "paymentRetry");
            if (inserted == 0) {
                continue;  // มี instance อื่นทำแล้ว — skip
            }
            payment.charge(p);   // ⭐ ทำหลัง insert (ถ้า charge fail → transaction rollback ทั้ง dedup row ด้วย)
        }
    }
}

กุญแจ:

  1. FOR UPDATE SKIP LOCKED — row-level lock (การล็อกทีละแถว ป้องกันไม่ให้ worker อื่นหยิบแถวที่กำลังประมวลผลอยู่) ใน Postgres → 2 worker หยิบ row ไม่ทับกัน
  2. INSERT ... ON CONFLICT หรือ unique key + catch — กัน duplicate ถึงแม้ ShedLock lock fail
  3. Transaction ครอบทั้ง insert + charge → atomic
  4. Idempotency key ควรรวม "retry attempt" → retry ครั้งหน้าใช้ key ใหม่ได้

Outbox Pattern กับ Batch

ถ้า batch ต้อง publish event (เช่น "report generated" → notify users) → ใช้ Outbox Pattern เหมือนกัน — insert outbox row ใน transaction เดียวกับ job result แล้วให้ relay process ส่ง event ออกทีหลัง วิธีนี้ทำให้ "บันทึก job result" กับ "ส่ง event" เป็น atomic กันโดยไม่ต้องพึ่ง distributed transaction (ดูรายละเอียดเพิ่มเติมใน บท 8 หัวข้อ Outbox Pattern — ไม่จำเป็นต้องอ่านก่อน)

5.2 Monitoring

batch job รันเงียบ ๆ เบื้องหลัง จึงต้อง monitor ให้รู้เมื่อมีปัญหา — Spring Batch เก็บ metadata ทุก run ในตาราง (BATCH_JOB_EXECUTION) ที่ query มาทำ dashboard ได้ (สถานะ, duration, read/write/skip count) หรือ export เป็น Micrometer metric:

Spring Batch tables (BATCH_JOB_EXECUTION, ฯลฯ) — query เป็น dashboard:

  • Last run status
  • Duration
  • Read/write/skip count

หรือ export metric ผ่าน Micrometer:

💡 Spring Batch 5 + Spring Boot 3.x ส่ง metric (spring.batch.job.*, spring.batch.step.*) เข้า Micrometer ให้อัตโนมัติเมื่อมี MeterRegistry อยู่ใน classpath — ไม่ต้องสร้าง @Bean เอง เพียงเพิ่ม dependency ต่อไปนี้:

xml
<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-actuator</artifactId>
</dependency>
<dependency>
    <groupId>io.micrometer</groupId>
    <artifactId>micrometer-registry-prometheus</artifactId>
</dependency>

แล้ว metric จะขึ้นมาที่ /actuator/metrics และ /actuator/prometheus เอง

5.3 Alert ที่:

metric เปล่า ๆ ไม่พอ ต้องตั้ง alert ที่จับสัญญาณผิดปกติของ batch โดยเฉพาะ — งานที่ไม่ run ตามเวลา (cron missed อันตรายเพราะเงียบ), duration ผิดปกติ, skip count สูง, หรือ status FAILED:

  • Job ไม่ run ตามเวลา (cron missed)
  • Job duration เกิน baseline 2x
  • Skip count สูงเกิน threshold
  • Job status = FAILED

5.4 Timeout

job ที่ค้าง (infinite loop, DB hang) จะกินทรัพยากรไม่ปล่อยและบล็อกรอบถัดไป — ควรมี timeout guard ที่ตรวจเวลาที่ใช้ไปแล้ว break ออกเมื่อเกิน threshold กันไม่ให้ job หนึ่งลากทั้งระบบ:

Job ที่ runaway = หาก infinite loop / DB hang:

java
@Scheduled(cron = "...")
public void job() {
    long start = System.currentTimeMillis();
    Iterator<X> it = ...;
    while (it.hasNext()) {
        if (System.currentTimeMillis() - start > 30 * 60 * 1000) {
            log.warn("timeout, abort");
            break;
        }
        process(it.next());
    }
}

5.5 Resource Isolation

Batch job ใหญ่ ๆ — แยก JVM / pod จาก online traffic

  • Online API: 4 pod × 2 CPU
  • Batch: 1 pod × 8 CPU (รัน ตี 2)

ตอนรัน batch — online ไม่ slow


Part 6: Lab

Lab 1: @Scheduled + ShedLock

ทำ scheduler ที่ลบ user ที่ unverified > 7 วัน

  • รัน 3 instance — verify ว่า lock ทำงาน (DB เห็น delete_unverified row หมุน lock)
  • วิธีรัน 3 instance บนเครื่องเดียว (เปิด 3 terminal คนละ port):
    bash
    ./mvnw spring-boot:run -Dspring-boot.run.arguments=--server.port=8080
    ./mvnw spring-boot:run -Dspring-boot.run.arguments=--server.port=8081
    ./mvnw spring-boot:run -Dspring-boot.run.arguments=--server.port=8082

Lab 2: Spring Batch ETL CSV → DB

ฝึก Spring Batch กับงานจริง — อ่าน CSV ขนาดใหญ่ (1M row), validate, insert DB โดยใช้ chunk + retry + skip invalid row จะเห็นว่า batch จัดการข้อมูลใหญ่และ error ได้ดีกว่าเขียน loop เอง (และ restart ได้ถ้าพังกลางทาง):

  1. อ่าน users.csv (1M row)
  2. Validate email
  3. Insert DB (skip duplicate)
  4. ใช้ chunk 1000, retry 3, skip บน invalid email

Lab 3: Quartz Dynamic Reminder

API ให้ user ตั้งเวลา reminder → schedule Quartz job → ส่ง email ตามเวลา


Part 7: Checkpoint

  1. fixedRate ต่างจาก fixedDelay ยังไง?
  2. ทำไม @Scheduled มีปัญหาตอน scale หลาย pod? แก้ยังไง?
  3. Quartz มี feature อะไรที่ @Scheduled ไม่มี?
  4. Cron field ของ Spring มีกี่ field?
  5. Spring Batch "chunk" ทำงานยังไง?
  6. Restart ใน Spring Batch ทำงานยังไง?
  7. Idempotent job สำคัญยังไง?
  8. Partition step ใช้ทำอะไร?
  9. เมื่อไหร่ใช้ K8s CronJob แทน @Scheduled?
  10. ทำไมต้องแยก batch JVM จาก online traffic?

Part 8: สรุปบทนี้

  • @Scheduled — ง่าย แต่ scale หลาย pod ต้องใช้ ShedLock
  • Quartz — persistent + cluster + dynamic schedule
  • Spring Batch — chunk + restart + retry + skip + parallel — เหมาะ ETL/report
  • K8s CronJob — แยก scheduling ออกจาก app
  • Idempotent = job รันซ้ำได้
  • Monitor: cron missed, duration, skip rate, FAIL status
  • Production: resource isolation, timeout, alert, job metadata dashboard

บทถัดไป — GraphQL (Spring for GraphQL)


← บทที่ 12 | สารบัญ | บทที่ 14: GraphQL