Skip to content

บทที่ 19 — Polyglot (พ็อล-ลี-ก็อต = "พูดหลายภาษา") Example End-to-End

← บทที่ 18: Deployment | สารบัญ

TL;DR: บทนี้คือ capstone — รวม 18 บทก่อนหน้าเป็นระบบจริง: Order (Java/Spring) + Inventory (Go) + Notification (Node/NestJS) + Recommendation (Python/FastAPI) + Kafka + RabbitMQ + Postgres + Redis + K8s + Istio + Argo CD + OTel + Jaeger. แสดง pattern เดียวกัน (Outbox + Inbox + idempotent) ใช้ได้ทุกภาษา. ข้ามได้ถ้า: อ่านบท 0-18 แล้วเห็นภาพ end-to-end อยู่แล้ว — แต่ถ้าอยากเห็นโค้ดจริง บทนี้คือ checkpoint รวมสุดท้าย

🗺️ Minimum Viable Path (ยังไม่คุ้น K8s/Istio/Argo CD): ไม่เป็นไร — ส่วนนั้น optional; เน้นอ่าน Part 3–6 (Order/Inventory/Notification/Recommendation flow ผ่าน Kafka) ก็เข้าใจ polyglot pattern ได้แล้ว. K8s + Istio + Argo CD (Part 7) อ่านเพิ่มเติมได้ภายหลังเมื่อผ่านบท 17-18 แล้ว

📖 Polyglot = poly (หลาย) + glot (ลิ้น/ภาษา) = "พูดหลายภาษา". ในบริบทนี้ = ระบบ microservices ที่ใช้ภาษาโปรแกรมหลายภาษาผสมกัน (Java + Go + Node + Python) แต่ละ service เลือกภาษาที่เหมาะกับงาน

💡 มือใหม่ที่เพิ่งจบ Java/Spring: ตัวอย่างใช้หลายภาษา — ไม่ต้องเข้าใจ syntax ของ Go/Python/Node ทุกบรรทัด. เน้นที่ flow + interaction ระหว่าง service (ใคร publish event อะไร, ใคร consume, pattern อะไร) — syntax เป็นรอง

🛠️ โค้ดที่ run ได้: lab จริงมี 4 service 4 ภาษาทำงานผ่าน Kafka (Java/Spring + Go + Node/TS + Python/FastAPI). เนื้อหาในบทเล่าครบ stack ใหญ่ (รวม Istio + Argo CD) — lab โฟกัสที่ "polyglot ผ่าน Kafka" 4 ภาษาที่เล่มสอน (ลิงก์ lab รออัพเดต — กำลังย้าย repo ไป public host ที่ stable; กำหนดประมาณ Q3 2026 บน GitHub Pages/GitLab public — รอติดตามในตอนเริ่มต้นของ README หลัก)

ทดลองก่อน lab พร้อม: ระหว่างรอ lab หลัก — ลองสร้าง Docker Compose ที่มี 4 service + Kafka (Bitnami/Apache image) + Postgres แยกต่อ service แล้ว copy code จาก Part 3-6 ในบทนี้ เป็นการฝึกที่ดีและเห็นผล polyglot pattern จริง. ดู docker-compose.yml ตัวอย่างได้ที่ README หลักเมื่อ lab พร้อม

บทสุดท้ายของหมวด Microservices — รวมทุกอย่างใน 18 บทก่อนหน้า เป็นระบบจริงที่ใช้:

  • Java (Spring Boot) — Order Service
  • Go — Inventory Service
  • Node.js (NestJS) — Notification Service
  • Python (FastAPI) — Recommendation Service
  • Kafka — event log
  • RabbitMQ — task queue
  • PostgreSQL — per-service DB
  • Redis — cache + Streams
  • Kubernetes + Istio + Argo CD — orchestration + mesh + GitOps
  • OpenTelemetry + Jaeger + Prometheus + Grafana — observability

ใช้เวลา 4-5 ชั่วโมง เพื่อ build + run


Part 1: Use Case — Mini E-commerce

text
User flow:
0. User คลิกปุ่ม "Order" ใน React app (frontend)
   → frontend ส่ง POST /orders มาที่ API Gateway → Order Service
1. POST /orders  → Order Service (Java)
2. Order Service publishes OrderPlaced → Kafka
3. Inventory Service (Go) consumes → reserve stock
4. Notification Service (Node) consumes → send email/SMS
5. Recommendation Service (Python) consumes → update model
6. Async after some time → Order confirmed

Part 2: Architecture


Part 3: Service 1 — Order Service (Java + Spring Boot)

3.1 Domain + Outbox

เริ่มที่ Order Service (Java) — สอง entity หลักคือ Order (ข้อมูลคำสั่งซื้อ) และ OutboxEvent (ตาราง outbox สำหรับ publish event อย่างปลอดภัย) ทั้งคู่อยู่ DB เดียวกันเพื่อให้บันทึกพร้อมกันใน transaction เดียวได้ (กัน dual-write):

java
@Entity
@Table(name = "orders")
public class Order {
    @Id UUID id;
    String customerId;
    BigDecimal total;
    String status;    // PLACED, CONFIRMED, CANCELLED
    Instant createdAt;
}

@Entity
@Table(name = "outbox")
public class OutboxEvent {
    @Id UUID id;
    String aggregateType;
    String aggregateId;
    String eventType;
    @Convert(converter = JsonbConverter.class)
    JsonNode payload;
    Instant createdAt;
    Instant sentAt;
}

3.2 Service

📌 Stack ของ Order Service: Java 21 + Spring Boot 3 + Spring Data JPA + PostgreSQL (ใช้ JSONB column สำหรับ payload event)

หัวใจของ business logic — method place สร้าง order แล้วบันทึก outbox event ("OrderPlaced") ใน @Transactional เดียวกัน ถ้า commit สำเร็จก็การันตีว่าทั้ง order และ event ถูกบันทึกพร้อมกัน (ตัว relay/CDC จะ publish event ทีหลัง):

java
@Service
public class OrderService {
    private final OrderRepo orderRepo;
    private final OutboxRepo outboxRepo;

    @Transactional
    public Order place(CreateOrderRequest req) {
        Order o = new Order();
        o.setId(UUID.randomUUID());
        o.setCustomerId(req.customerId());
        o.setTotal(req.total());
        o.setStatus("PLACED");
        orderRepo.save(o);

        OutboxEvent e = new OutboxEvent();
        e.setAggregateType("Order");
        e.setAggregateId(o.getId().toString());
        e.setEventType("OrderPlaced");
        // ⚠️ payload ต้องมีทุก field ที่ consumer ต้องใช้ — ฝั่ง Inventory (Go)
        //    ต้องการ eventId (dedup ใน inbox), productId, qty เพื่อจองของ
        //    ถ้าลืม field เหล่านี้ Go จะได้ค่า zero-value (qty=0) แล้วจองของผิดเงียบ ๆ
        e.setPayload(mapper.valueToTree(Map.of(
            "eventId", UUID.randomUUID(),   // UUID ต่อ event ใช้ dedup ฝั่ง consumer
            "orderId", o.getId(),
            "customerId", o.getCustomerId(),
            "productId", req.productId(),
            "qty", req.qty(),
            "total", o.getTotal()
        )));
        outboxRepo.save(e);

        return o;
    }
}

📌 CreateOrderRequest ต้องมี field ครบ: record นี้ต้องมี customerId(), productId(), qty(), total() เพราะ place() เอา productId/qty ไปใส่ใน payload ให้ Inventory (Go) ใช้จองของ — field name ใน JSON payload ต้องตรงกับ struct tag ฝั่ง Go เป๊ะ ๆ (eventId, orderId, customerId, productId, qty) ไม่งั้น consumer จะ parse ไม่เจอแล้วได้ค่า default แทน. ตอนนี้เราจับคู่ field ด้วยมือ (ad-hoc JSON) — Part 7.5 จะแสดงวิธีที่ดีกว่าด้วย schema กลาง (Protobuf/Avro) ที่ generate struct ทั้งสองฝั่งจาก contract เดียว

3.3 Controller

ชั้น HTTP ที่รับ request — รับ Idempotency-Key เพื่อกัน client retry สร้าง order ซ้ำ (ถ้า key ซ้ำก็คืนผลเดิม) และ X-User-Id ที่ gateway ฉีดมา แล้วเรียก service ทำงานจริง รวมหลายเทคนิคจากบทก่อน ๆ เข้าด้วยกัน:

java
@RestController
@RequestMapping("/orders")
public class OrderController {

    @PostMapping
    public ResponseEntity<OrderResponse> place(
            @RequestHeader("Idempotency-Key") String key,
            @RequestHeader("X-User-Id") String userId,
            @RequestBody @Valid CreateOrderRequest req) {
        
        Optional<Order> existing = idempotencyStore.get(key);
        if (existing.isPresent()) return ResponseEntity.ok(toResp(existing.get()));

        Order o = orderService.place(req);
        idempotencyStore.save(key, o);
        return ResponseEntity.status(HttpStatus.CREATED).body(toResp(o));
    }
}

3.4 Consume Event

Order Service ไม่ได้แค่ส่ง event แต่ยังฟังกลับด้วย — เมื่อ Inventory ตอบ StockReserved/StockReservationFailed ก็อัปเดตสถานะ order ตาม สังเกตการใช้ inbox (tryInsert) ทำ idempotent กัน event ซ้ำก่อน process

📖 Topic inventory.events มาจากไหน? สรุปสั้น: Inventory Service เขียน outbox record ลง DB → Debezium อ่าน WAL แล้ว publish เป็น Kafka event → ได้ topic inventory.events (รายละเอียด Outbox + Debezium ดู บทที่ 12 §CDC)

สำหรับผู้ที่ต้องการ config detail: Debezium ใช้ SMT (io.debezium.transforms.outbox.EventRouter) อ่าน column aggregate_type = "Inventory" แล้ว route topic ตาม pattern route.topic.replacement = ${routedByValue}.events → topic inventory.events — ไม่ต้องเข้าใจทุก field ตอนนี้ ดูบทที่ 12 เมื่อต้องการ configure จริง

📖 Polymorphic deserialization คืออะไร? = การแปลง JSON กลับเป็น object ที่ "มีได้หลายชนิด (subtype)" — ในที่นี้ event ที่มาจาก inventory.events อาจเป็นได้ทั้ง StockReserved (จองสำเร็จ) หรือ StockReservationFailed (จองไม่ได้). JSON เปล่า ๆ ไม่ได้บอกว่าตัวเองเป็นชนิดไหน Jackson จึงต้องมี "ป้ายบอก" (hint) ว่าให้เลือก subtype ไหนจาก field ตัวไหนใน JSON

⚠️ Polymorphic JSON deser — ต้องเซตให้ Jackson รู้ subtype: ถ้า payload เป็น JSON ปกติ Jackson ไม่รู้ จะแปลงเป็น StockReserved หรือ StockReservationFailed ตัวอย่างเลือกวิธีง่ายสุด: ใส่ @JsonTypeInfo + @JsonSubTypes บน sealed interface (Java 21+) — ทางเลือกอื่นคือ Custom JsonDeserializer ใน ConsumerFactory หรือใช้ Avro/Protobuf + Schema Registry (แนะนำสำหรับ prod)

java
// 1) ประกาศ sealed interface + ลงทะเบียน subtype ให้ Jackson
@JsonTypeInfo(use = JsonTypeInfo.Id.NAME, property = "_type")
@JsonSubTypes({
    @JsonSubTypes.Type(value = StockReserved.class,           name = "StockReserved"),
    @JsonSubTypes.Type(value = StockReservationFailed.class,  name = "StockReservationFailed")
})
public sealed interface InventoryEvent
        permits StockReserved, StockReservationFailed {
    UUID eventId();
    String orderId();
}

public record StockReserved(UUID eventId, String orderId, String productId, int qty)
        implements InventoryEvent {}
public record StockReservationFailed(UUID eventId, String orderId, String reason)
        implements InventoryEvent {}

// 2) Listener — groupId ใช้ชื่อเฉพาะ (กัน partition assignment ชนกับ listener อื่นใน service เดียวกัน)
@KafkaListener(topics = "inventory.events", groupId = "order-service-inventory-listener")
@Transactional
public void onInventoryEvent(InventoryEvent e) {
    if (e instanceof StockReserved sr) {
        if (!inboxRepo.tryInsert(sr.eventId())) return;  // idempotent
        orderRepo.updateStatus(sr.orderId(), "CONFIRMED");
    } else if (e instanceof StockReservationFailed srf) {
        if (!inboxRepo.tryInsert(srf.eventId())) return;
        orderRepo.updateStatus(srf.orderId(), "CANCELLED");
    }
}

💡 Producer ฝั่ง Inventory ต้องใส่ field _type ใน JSON ให้ตรงกับชื่อใน @JsonSubTypes (เช่น "_type":"StockReserved") — ถ้าใช้ Debezium outbox SMT ให้ใส่ _type ลงใน payload JSON ตอน write outbox


Part 4: Service 2 — Inventory Service (Go)

4.1 Schema

ต่อมาคือ Inventory Service เขียนด้วย Go (แสดงความเป็น polyglot) — schema มีตาราง stock (ของในคลัง) และ inbox สำหรับทำ idempotent consumer เหมือนฝั่ง Java แต่คนละภาษา ตอกย้ำว่า pattern เดียวกันใช้ได้ทุกภาษา:

sql
CREATE TABLE stock (
    product_id TEXT PRIMARY KEY,
    on_hand INT NOT NULL,
    reserved INT NOT NULL DEFAULT 0
);

CREATE TABLE inbox (
    event_id UUID PRIMARY KEY,
    received_at TIMESTAMPTZ DEFAULT now()
);

CREATE TABLE outbox (
    id UUID PRIMARY KEY,
    event_type TEXT,
    payload JSONB,
    created_at TIMESTAMPTZ DEFAULT now(),
    sent_at TIMESTAMPTZ
);

4.2 Code

นี่คือ consumer ของ Inventory แบบเต็มใน Go ทำ 4 ขั้นในหนึ่ง transaction เดียว:

  1. อ่าน event จาก topic orders.events
  2. กัน event ซ้ำ (idempotent) ด้วยตาราง inbox
  3. จองของUPDATE ... WHERE on_hand - reserved >= qty
  4. เขียน outbox event ผลลัพธ์StockReserved (สำเร็จ) หรือ StockReservationFailed (ของหมด)

สังเกตว่าใช้ pattern เดียวกับฝั่ง Java เป๊ะ แค่เปลี่ยนภาษา:

go
package main

import (
    "context"
    "database/sql"
    "encoding/json"
    "log"
    "github.com/segmentio/kafka-go"
    _ "github.com/jackc/pgx/v5/stdlib"
    "github.com/google/uuid"
)

type OrderPlaced struct {
    EventId    string  `json:"eventId"`   // UUID ต่อ event (ใช้ dedup ใน inbox)
    OrderId    string  `json:"orderId"`
    CustomerId string  `json:"customerId"`
    ProductId  string  `json:"productId"`
    Qty        int     `json:"qty"`
}

func main() {
    db, _ := sql.Open("pgx", "postgres://...")

    // CommitInterval: 0 → manual commit เท่านั้น (กัน auto-commit ก่อน tx จริงสำเร็จ)
    reader := kafka.NewReader(kafka.ReaderConfig{
        Brokers:        []string{"kafka:9092"},
        Topic:          "orders.events",
        GroupID:        "inventory-service",
        CommitInterval: 0,
    })

    for {
        ctx := context.Background()
        msg, err := reader.FetchMessage(ctx)   // FetchMessage = manual commit mode
        if err != nil { log.Println(err); continue }

        var event OrderPlaced
        if err := json.Unmarshal(msg.Value, &event); err != nil {
            // bad payload → skip + commit เพื่อไม่ให้ stuck (หรือส่ง DLQ)
            _ = reader.CommitMessages(ctx, msg)
            continue
        }

        if err := handleOrderPlaced(db, event); err != nil {
            log.Printf("error: %v (will retry — not committing offset)", err)
            continue   // ไม่ commit → message จะถูก re-read หลัง rebalance/restart
        }

        // commit เฉพาะเมื่อ transaction สำเร็จ (at-least-once + outbox + inbox = effectively exactly-once)
        if err := reader.CommitMessages(ctx, msg); err != nil {
            log.Printf("commit failed: %v", err)
        }
    }
}

func handleOrderPlaced(db *sql.DB, e OrderPlaced) error {
    tx, err := db.Begin()
    if err != nil { return err }
    defer tx.Rollback()

    // Idempotent dedup ที่ระดับ "event" — ใช้ EventId (UUID ต่อ event)
    // ห้ามใช้ OrderId เพราะ order เดียวกันอาจมีหลาย event (Placed/Updated/...)
    result, err := tx.Exec("INSERT INTO inbox(event_id) VALUES($1) ON CONFLICT DO NOTHING", e.EventId)
    if err != nil { return err }
    // ⚠️ ต้องตรวจ RowsAffected: ON CONFLICT DO NOTHING คืน 0 rows ถ้า event_id ซ้ำ
    // ถ้าไม่ตรวจ → duplicate event จะยังถูก process ต่อ (inbox pattern พัง)
    // อย่ากลืน error ของ RowsAffected() ด้วย `_` — ถ้า driver คืน error แล้วเราเผลอมองเป็น 0
    //   จะเข้าใจผิดว่า "event ซ้ำ" แล้วข้าม event จริงไปเงียบ ๆ
    count, err := result.RowsAffected()
    if err != nil { return err }
    if count == 0 {
        return tx.Commit() // duplicate — commit no-op แล้วออก
    }

    // Try reserve
    res, err := tx.Exec(`
        UPDATE stock SET reserved = reserved + $1
        WHERE product_id = $2 AND on_hand - reserved >= $1
    `, e.Qty, e.ProductId)
    rows, _ := res.RowsAffected()

    var outboxEvent map[string]interface{}
    if rows > 0 {
        outboxEvent = map[string]interface{}{
            "type": "StockReserved",
            "orderId": e.OrderId,
            "productId": e.ProductId,
            "qty": e.Qty,
        }
    } else {
        outboxEvent = map[string]interface{}{
            "type": "StockReservationFailed",
            "orderId": e.OrderId,
            "reason": "out_of_stock",
        }
    }
    
    payload, _ := json.Marshal(outboxEvent)
    _, err = tx.Exec(`
        INSERT INTO outbox(id, event_type, payload)
        VALUES($1, $2, $3)
    `, uuid.New(), outboxEvent["type"], payload)
    if err != nil { return err }

    return tx.Commit()
}

Debezium captures outbox → publishes to inventory.events

⚠️ loop นี้ยังขาด backoff: ตอน handler error โค้ดทำ continue ทันทีเพื่อ re-read message — ถ้า downstream ล่มค้างนาน (เช่น DB down) มันจะวน retry ถี่ ๆ ไม่มีหยุด (tight loop) กระหน่ำทั้ง DB และ Kafka broker. ในของจริงควรใส่ exponential backoff + jitter (ดูบท 13) ก่อน retry รอบถัดไป หรือกำหนดจำนวน retry สูงสุดแล้วโยนเข้า DLQ — เหมือนที่ฝั่ง Notification worker ทำ


Part 5: Service 3 — Notification (Node.js + NestJS)

5.1 Consumer

Service ที่สามคือ Notification (Node.js + NestJS) — ภาษาที่สามในระบบเดียว consumer ฟัง order.events แล้วเมื่อเจอ OrderPlaced ก็โยนงานส่งอีเมลเข้า RabbitMQ task queue (แยกงานหนัก/ช้าออกจาก hot path):

typescript
import { Controller } from '@nestjs/common';
import { EventPattern, Payload } from '@nestjs/microservices';

@Controller()
export class NotificationController {

  constructor(
    private email: EmailService,
    private rabbitMq: RabbitMqService,
  ) {}

  @EventPattern('order.events')
  async onOrderEvent(@Payload() data: any) {
    if (data.type === 'OrderPlaced') {
      // Send to RabbitMQ task queue (delay job)
      await this.rabbitMq.publish('tasks.email', {
        template: 'order-confirmation',
        to: data.customerEmail,
        data: { orderId: data.orderId, total: data.total }
      });
    }
  }
}

5.2 RabbitMQ Worker

worker ที่ดึงงานจาก task queue มาส่งอีเมลจริง — ack เมื่อสำเร็จ, ถ้า fail ก็ retry (nack requeue) จนครบ 3 ครั้งแล้วโยนเข้า DLQ (reject) เป็นตัวอย่าง competing-consumer + DLQ pattern จากบท 06/07 ในงานจริง:

typescript
@Controller()
export class EmailWorker {

  @MessagePattern('tasks.email')
  async sendEmail(@Payload() job: EmailJob, @Ctx() ctx: RmqContext) {
    const channel = ctx.getChannelRef();
    const msg = ctx.getMessage();
    
    try {
      await this.emailService.send(job);
      channel.ack(msg);
    } catch (e) {
      // retry up to 3 times → DLQ
      // ⚠️ ห้ามใช้ nack(requeue=true) แล้วคาดว่าจะนับ retry — header ไม่ถูก increment
      //    → จะ requeue วนไม่รู้จบ ไม่เข้า DLQ
      //    วิธีถูก: republish ไปที่ delay queue พร้อม header x-retry ใหม่ แล้ว ack ของเดิม
      const retryCount = (msg.properties.headers?.['x-retry'] as number) ?? 0;
      if (retryCount >= 3) {
        channel.reject(msg, false);   // → DLQ (queue ต้องตั้ง x-dead-letter-exchange)
      } else {
        const newHeaders = {
          ...msg.properties.headers,
          'x-retry': retryCount + 1,
        };
        // ส่งใหม่เข้า "tasks.email.retry" (delay queue) — TTL หมดแล้วเด้งกลับ tasks.email
        channel.publish(
          '',                          // default exchange
          'tasks.email.retry',         // routing key = delay queue name
          msg.content,
          { ...msg.properties, headers: newHeaders, persistent: true }
        );
        channel.ack(msg);              // ack ของเดิม (กัน double-process)
      }
    }
  }
}

📖 Delay queue pattern: ดูบทที่ 07 Part 10.2 — ตั้ง queue tasks.email.retry ที่ไม่มี consumer + x-message-ttl: 30000 + x-dead-letter-exchange: "" + x-dead-letter-routing-key: tasks.email → message รอ 30s แล้วเด้งกลับ queue หลัก เป็น exponential backoff retry ได้


Part 6: Service 4 — Recommendation (Python + FastAPI)

⚠️ กฎสำคัญ 2 ข้อสำหรับ FastAPI 2026:

  1. @app.on_event("startup"/"shutdown") deprecated ตั้งแต่ FastAPI 0.93 (2023) → ใช้ lifespan async context manager แทน
  2. ห้ามเรียก sync API ใน async functionconfluent_kafka.Consumer.poll() เป็น blocking sync → ถ้าใส่ใน async def จะ block event loop ทั้ง FastAPI → request อื่นค้างหมด วิธีถูก: ใช้ aiokafka (AsyncIO-native) หรือ wrap ด้วย await asyncio.to_thread(consumer.poll, 1.0)
python
from contextlib import asynccontextmanager
from fastapi import FastAPI
from aiokafka import AIOKafkaConsumer  # ✅ AsyncIO-native (แนะนำ)
import asyncio, json

# --- 1) lifespan แทน on_event("startup") ---
@asynccontextmanager
async def lifespan(app: FastAPI):
    # startup
    task = asyncio.create_task(consume_events())
    yield
    # shutdown
    task.cancel()
    try:
        await task
    except asyncio.CancelledError:
        pass

app = FastAPI(lifespan=lifespan)

# --- 2) consumer ใช้ aiokafka (ไม่บล็อก event loop) ---
async def consume_events():
    consumer = AIOKafkaConsumer(
        'orders.events',
        bootstrap_servers='kafka:9092',
        group_id='recommendation-service',
        auto_offset_reset='earliest',
        enable_auto_commit=False,   # commit เองหลัง process สำเร็จ
    )
    await consumer.start()
    try:
        async for msg in consumer:
            try:
                event = json.loads(msg.value)
                if event.get('type') == 'OrderPlaced':
                    await update_model(event['customerId'], event['productId'])
                await consumer.commit()
            except Exception as e:
                # log แล้วไม่ commit → message จะถูก re-read
                print(f"error: {e}")
    finally:
        await consumer.stop()

@app.get("/recommend/{user_id}")
async def recommend(user_id: str):
    return {"items": await model.recommend(user_id)}

💡 ทางเลือก (ยังใช้ confluent-kafka sync ได้): ห่อด้วย asyncio.to_thread:

python
msg = await asyncio.to_thread(consumer.poll, 1.0)

จะ run poll() ใน thread pool — ไม่บล็อก event loop แต่จ่าย overhead thread switching


Part 7: K8s Deployment

7.1 Order Service

manifest deploy ของ Order Service ที่รวมทุกอย่างที่ดี — 3 replica, inject Istio sidecar, ดึง secret จาก K8s Secret, ตั้ง liveness/readiness probe แยกกัน, resource request/limit และ HPA scale ตาม CPU เป็น template ที่ production จริงควรมี:

yaml
apiVersion: apps/v1
kind: Deployment
metadata: { name: order-service, namespace: production }
spec:
  replicas: 3
  selector: { matchLabels: { app: order } }
  template:
    metadata:
      labels: { app: order }
      annotations: { sidecar.istio.io/inject: "true" }
    spec:
      serviceAccountName: order-sa
      containers:
        - name: app
          image: ghcr.io/acme/order-service:abc1234
          env:
            - name: SPRING_DATASOURCE_PASSWORD
              valueFrom: { secretKeyRef: { name: db-secret, key: password } }
            - name: SPRING_KAFKA_BOOTSTRAP_SERVERS
              value: "kafka:9092"
            - name: OTEL_EXPORTER_OTLP_ENDPOINT
              value: "http://otel-collector:4317"
          resources:
            requests: { cpu: 200m, memory: 512Mi }
            limits: { cpu: 1000m, memory: 1Gi }
          readinessProbe:
            httpGet: { path: /actuator/health/readiness, port: 8080 }
            initialDelaySeconds: 10
            periodSeconds: 5
          livenessProbe:
            httpGet: { path: /actuator/health/liveness, port: 8080 }
            initialDelaySeconds: 30
            periodSeconds: 30
---
apiVersion: v1
kind: Service
metadata: { name: order, namespace: production }
spec:
  selector: { app: order }
  ports: [{ port: 80, targetPort: 8080 }]
---
apiVersion: autoscaling/v2
kind: HorizontalPodAutoscaler
metadata: { name: order-hpa }
spec:
  scaleTargetRef: { apiVersion: apps/v1, kind: Deployment, name: order-service }
  minReplicas: 3
  maxReplicas: 20
  metrics:
    - type: Resource
      resource: { name: cpu, target: { type: Utilization, averageUtilization: 70 } }

7.2 Istio mTLS + Authorization

เปิด zero-trust ด้วย Istio — PeerAuthentication mode STRICT บังคับ mTLS ทุก traffic เข้า service และ AuthorizationPolicy ระบุว่าเฉพาะ gateway/bff เท่านั้นเรียก order ได้ (ตาม ServiceAccount) นำแนวคิดบท 16 มาใช้จริง:

yaml
apiVersion: security.istio.io/v1
kind: PeerAuthentication
metadata: { name: default, namespace: production }
spec: { mtls: { mode: STRICT } }
---
apiVersion: security.istio.io/v1
kind: AuthorizationPolicy
metadata: { name: order-svc, namespace: production }
spec:
  selector: { matchLabels: { app: order } }
  rules:
    - from:
        - source:
            principals:
              - "cluster.local/ns/production/sa/gateway"
              - "cluster.local/ns/production/sa/bff-mobile"
      to:
        - operation: { methods: [GET, POST], paths: [/orders/*] }

7.3 Argo Rollout (Canary)

deploy เวอร์ชันใหม่แบบ canary ด้วย Argo Rollouts — กำหนด step ค่อย ๆ เพิ่ม weight (10% → 25% → 50% → 100%) พร้อม pause และ analysis ตรวจ error rate ระหว่างทาง ถ้าผิดปกติจะ rollback อัตโนมัติ ปิดท้ายด้วยการ deploy ที่ปลอดภัย:

yaml
apiVersion: argoproj.io/v1alpha1
kind: Rollout
metadata: { name: order-service }
spec:
  replicas: 10
  strategy:
    canary:
      steps:
        - setWeight: 10
        - pause: { duration: 5m }
        - analysis: { templates: [{ templateName: error-rate-check }] }
        - setWeight: 25
        - pause: { duration: 10m }
        - setWeight: 50
        - pause: { duration: 10m }
        - setWeight: 100

Part 7.5: Shared Contracts — gRPC + Protobuf ข้ามภาษา

ระบบ polyglot 4 ภาษาจะอยู่ด้วยกันได้ต้องมี contract กลาง. JSON ad-hoc = ของจริงพังเร็ว — ใช้ schema ภาษากลางแล้ว generate code ทุกภาษา

ตัวเลือกมาตรฐาน 2026:

Toolใช้ตอนไหน
Protocol Buffers v3 + gRPCsync RPC ข้าม service (low-latency, binary) — optional keyword stable ตั้งแต่ proto3.15 (2021) ใช้ได้เลย (ใน proto3 ปกติทุก field มีค่า default เสมอ ไม่มี null — keyword optional ทำให้เช็คได้ว่า field นั้นถูกส่งมาจริงหรือไม่)
🆕 Connect (Buf)gRPC-compatible แต่รัน HTTP/1.1 + HTTP/2 + HTTP/3 ด้วย JSON หรือ binary — debug ผ่าน curl ได้ ฝั่ง browser ไม่ต้อง grpc-web proxy
Avro + Confluent Schema Registryevent payload ใน Kafka (compact + schema evolution rules ดี)
JSON Schema + AsyncAPIdoc + lint event payload (ไม่ใช่ wire format)

Buf CLI + BSR (Buf Schema Registry) — toolchain มาตรฐาน 2026:

  • buf lint — บังคับ style guide (no breaking change, naming convention)
  • buf breaking — ตรวจ breaking change ใน CI ก่อน merge
  • buf generate — gen code ทุกภาษา (Java/Go/Node/Python/Kotlin/Rust) จาก buf.gen.yaml
  • BSR = Schema Registry สำหรับ Protobuf (เทียบเท่า Confluent Schema Registry แต่สำหรับ proto) — host schema, version, push/pull แทน vendoring .proto ลง repo ทุก service

ตัวอย่าง inventory.proto (shared):

proto
syntax = "proto3";
package acme.inventory.v1;

service InventoryService {
  rpc ReserveStock(ReserveStockRequest) returns (ReserveStockResponse);
}

message ReserveStockRequest {
  string order_id    = 1;
  string product_id  = 2;
  int32  qty         = 3;
  optional string idempotency_key = 4;   // proto3.15+ stable
}

message ReserveStockResponse {
  bool   reserved = 1;
  string reason   = 2;
}

buf generate ออก stub Java/Go/Node/Python ครบ — Order Service (Java) เรียก Inventory Service (Go) ผ่าน gRPC โดยไม่มีใครต้องเดา field name


Part 8: Observability

8.1 OpenTelemetry

แม้ service จะเขียนคนละภาษา (Java/Go/Node/Python) ทุกตัวส่ง telemetry ผ่าน OTel ไป Collector กลาง ทำให้ได้ trace ที่ร้อยเรียงข้ามทุก service เป็นเส้นเดียว (ตามตัวอย่าง trace flow ด้านล่าง) — นี่คือพลังของมาตรฐานกลาง:

ทุก service → OTel SDK → Collector → Jaeger + Tempo + Prometheus

Trace flow:

text
gateway.handle (50ms)
  └── order.create (30ms)
       ├── db.insert (5ms)
       ├── outbox.insert (3ms)
       └── kafka.send via debezium (async, after commit)
              └── inventory.consume (10ms)
                   ├── db.tryReserve (4ms)
                   └── outbox.insert (2ms)

8.2 Metrics

นอกจาก trace แต่ละ service ต้องเก็บ metric เชิงตัวเลขด้วย — ทั้ง metric ทั่วไป (request rate, latency P95, error rate) และเฉพาะ event-driven (Kafka lag, outbox unsent count) ที่บอกสุขภาพของ async flow ได้ดี:

Per service:

  • HTTP request rate, latency P95, error rate
  • Kafka lag
  • Outbox unsent count
  • DB pool wait time

8.3 Alerts

metric ต้องมี alert ที่จับสัญญาณก่อนผู้ใช้เจอปัญหา — นอกจาก 5xx/latency ปกติ ระบบ event-driven ต้องเฝ้า lag และ outbox unsent (ถ้าค้างเยอะ = relay/CDC ติด) บวก SLO burn rate ที่เตือนเมื่อ error budget หมดเร็วผิดปกติ:

text
- HTTP 5xx rate > 1% (5 min)
- Kafka consumer lag > 10000
- Outbox unsent > 1000 (something blocked)
- DB pool wait > 100ms
- Pod restart > 3 in 10min
- SLO burn rate alert

Part 9: End-to-end Test

typescript
// Cypress / Playwright
test('place order end-to-end', async () => {
  // 1. Create order via API
  const orderRes = await api.post('/orders', {
    headers: { 'Idempotency-Key': uuid(), 'X-User-Id': testUserId },
    body: { productId, qty: 2, total: 100 }
  });
  expect(orderRes.status).toBe(201);
  const orderId = orderRes.body.id;

  // 2. Wait for status confirmed (poll)
  await waitFor(async () => {
    const r = await api.get(`/orders/${orderId}`);
    return r.body.status === 'CONFIRMED';
  }, { timeout: 30_000 });

  // 3. Check inventory decremented
  const inv = await api.get(`/inventory/${productId}`);
  expect(inv.body.reserved).toBeGreaterThan(0);

  // 4. Check email sent (via mailtrap test)
  const emails = await mailtrap.getEmails(testUserEmail);
  expect(emails.some(e => e.subject.includes('Order confirmation'))).toBe(true);
});

Part 10: ทดสอบ Failure Scenarios

10.1 Kill Pod (Chaos)

ทดสอบว่าระบบทนต่อ pod ล่ม — จงใจ kill pod ของ inventory แล้วยืนยันว่า order ยังทำงานต่อได้ เพราะ consumer ที่เหลือ rejoin group และ reprocess message ที่ค้าง (อาศัย at-least-once + idempotent ที่วางไว้):

bash
kubectl delete pod -l app=inventory --random

→ Order ยัง process ได้ (consumer rejoin group, reprocess message)

10.2 Slow Downstream

ทดสอบว่าระบบรับมือ downstream ช้าได้ — ใช้ Istio fault injection จงใจหน่วง 5 วินาที แล้วดูว่า circuit breaker ของ Order Service เปิด (fail fast + fallback) แทนที่จะรอจน thread หมดและล้มเป็นโดมิโน (ตามบท 13):

bash
istioctl ... fault delay 5s 100%

→ Order Service circuit breaker open → fail fast + fallback

10.3 Kafka Down

ทดสอบว่า broker ล่มแล้ว event ไม่หาย — เพราะใช้ Outbox pattern ตอน Kafka ล่ม event จะสะสมในตาราง outbox (ไม่หาย) แล้ว relay/CDC จะส่งให้ครบเมื่อ Kafka กลับมา นี่คือเหตุผลที่ Outbox สำคัญ:

→ Outbox accumulates → no event loss

10.4 DB Down

→ Liveness still alive, readiness false → no traffic → no data loss


Part 11: สิ่งที่ได้เรียนใน 19 บท

จบหมวด Microservices คุณรู้:

  • เมื่อไหร่ใช้ / ไม่ใช้ (00)
  • ตัด service ด้วย DDD (01)
  • สื่อสาร sync (REST/gRPC) + async (broker) (02-03)
  • Gateway + BFF (04)
  • Discovery + Config (05)
  • Broker fundamentals + RabbitMQ + Kafka เข้ม (06-08)
  • ตัวเลือก broker อื่น (09-10)
  • Stream processing (11)
  • Saga + Outbox + Inbox + CDC (12) ← หัวใจ event-driven
  • Resilience (13)
  • Data management + CQRS + ES (14)
  • Distributed Tracing (15)
  • Zero Trust Security (16)
  • Service Mesh (17)
  • Deployment + GitOps (18)
  • Polyglot real system (19 — บทนี้)

Part 12: ต่อไปจะอ่านอะไร

หลังจบเล่มนี้ — เรียนต่อ:

  • Kubernetes ลึก (devops/04)
  • Observability ลึก (observability/)
  • Distributed Systems Theory (distributed-systems/)
  • System Design (system-design/)
  • Temporal / Camunda สำหรับ orchestration
  • eBPF / Cilium สำหรับ next-gen networking
  • Wasm in microservices (emerging)

Part 12.5: 📋 Production Readiness Checklist

ก่อน polyglot system นี้ขึ้น production จริง — ทบทวนทุกข้อ:

Architecture

text
☐ DDD bounded context ชัดเจน ทุก service
☐ DB per service (ไม่มี shared DB)
☐ Async event-driven สำหรับ business flow (Kafka)
☐ Sync เฉพาะ user-facing read (REST/gRPC + BFF)
☐ Idempotency-Key สำหรับทุก mutate API
☐ Outbox pattern + CDC (Debezium) สำหรับ DB → Kafka
☐ Inbox table สำหรับ consumer dedup
☐ Saga (Temporal/orchestration) สำหรับ multi-service flow

Reliability

text
☐ Timeout ทุก outbound call (HTTP, gRPC, DB, broker)
☐ Retry: idempotent only + exponential + jitter + budget
☐ Circuit breaker per downstream
☐ Bulkhead per critical dependency
☐ Hedged request สำหรับ read-only critical path
☐ Load shedding + priority queue
☐ Graceful shutdown (terminationGracePeriodSeconds: 60+)
☐ Chaos test quarterly (kill pod, latency inject, network partition)
☐ Capacity plan peak × 3

Observability (4 pillars)

text
☐ OTel SDK / Java agent ทุก service
☐ W3C traceparent + baggage propagation (HTTP + Kafka)
☐ OTel Collector + tail sampling
☐ Metrics: RED (Rate/Error/Duration), USE (Util/Sat/Errors)
☐ Logs: JSON structured + trace_id + correlation
☐ Profiling: pprof / async-profiler (OTel profiling)
☐ SLO + Error budget per service
☐ Alert runbook for top 10 alerts

Security

text
☐ mTLS ทุก service-to-service (service mesh)
☐ Workload identity (SPIFFE/SPIRE หรือ K8s SA)
☐ JWT validation ทุก service (Zero Trust)
☐ AuthorizationPolicy + OPA/OpenFGA
☐ Broker auth: SASL/SCRAM + ACL per topic
☐ Secret via External Secrets Operator + Vault
☐ NetworkPolicy default-deny
☐ Image: Cosign signed + Trivy scanned + SBOM
☐ Annual pentest + bug bounty

Deployment

text
☐ GitOps (Argo CD / Flux)
☐ Canary deploy (Argo Rollouts / Flagger)
☐ Auto-rollback on error rate spike
☐ Feature flag (OpenFeature) for risky feature
☐ DB migration: Expand-Contract (backward compat)
☐ Image digest, never `latest`
☐ PodDisruptionBudget set
☐ HPA configured (CPU + custom metric)
☐ Multi-AZ + multi-region (critical service)

Data

text
☐ Schema Registry (Avro/Protobuf)
☐ Outbox cleanup job (delete sent > 7 day)
☐ Debezium HA (2+ replica, monitor WAL lag)
☐ Backup per service + tested restore (quarterly)
☐ CDC → Lakehouse (Iceberg/Delta) for reporting
☐ PII tokenization + encryption at rest
☐ Data retention policy per table

Part 12.6: Production War Stories จาก End-to-End Polyglot System

🔥 บริษัท SaaS (2023): Avro schema fork ทำ Python consumer พัง 2 วัน

Java team update schema (add field optional) + register → Python consumer ใช้ schema เก่า → parse field ใหม่ผิด → crash loop

Lesson:

  • Schema Registry compatibility mode ต้องเป็น BACKWARD หรือ FULL (default)
  • Test consumer ทุกภาษาก่อน register schema
  • Auto-test ใน CI: register schema → run consumer ทุกภาษา → verify

🔥 บริษัท fintech (2024): Go consumer Kafka client เก่า → ไม่ support KIP-848

ทีม Go ใช้ confluent-kafka-go เก่า (0.5+) → KIP-848 ไม่ support → join consumer group ไม่ได้หลัง broker upgrade Kafka 4.0

Lesson:

  • Polyglot = ต้อง track client library version ของทุกภาษา
  • Test cluster upgrade กับ client ทุกภาษา/version ก่อน prod
  • Maintain "compatibility matrix" document

🔥 startup (2022): Istio sidecar ทำ Node.js (event-loop) latency spike 5x

Node.js single-thread + sidecar overhead → P99 spike ตอน sidecar GC

Lesson:

  • Migrate ไป Ambient Mesh (ztunnel + waypoint) — sidecar overhead หาย
  • หรือใช้ Cilium eBPF mesh สำหรับ Node.js-heavy workload
  • Benchmark mesh overhead per service type ก่อน enable

Part 13: Checkpoint

  1. End-to-end flow ของ order ตั้งแต่ user click จนเห็น email confirmation?
  2. Outbox + Inbox pattern ทำงานยังไงร่วมกันใน example นี้?
  3. ทำไม Order Service ใช้ Java, Inventory ใช้ Go, Notification ใช้ Node, Recommend ใช้ Python?
  4. Istio ทำอะไรใน example นี้?
  5. ถ้า Kafka down ระบบยังทำงานได้ไหม?
  6. ถ้า Inventory Service ตาย — order จะเกิดอะไร?
  7. Trace 1 request วิ่งผ่านกี่ service?
  8. Canary deploy ของ Order Service ทำงานยังไง?
  9. ทำไม Notification ใช้ทั้ง Kafka + RabbitMQ?
  10. ระบบนี้ scale แต่ละ service แยกได้ยังไง?

Part 14: สรุปบทนี้ + หมวด

  • Polyglot microservices = Java + Go + Node + Python ทำงานร่วมกันผ่าน contract (Avro/Protobuf)
  • Kafka = event log + audit
  • RabbitMQ = task queue
  • Outbox + Inbox + CDC (Debezium) = reliable messaging
  • Istio = mTLS + auth + traffic + observability
  • OpenTelemetry = unified observability
  • Argo CD + Rollouts = GitOps + progressive delivery
  • K8s + HPA = auto-scale
  • Resilience patterns + chaos test = production-grade

ขอแสดงความยินดี — คุณจบหมวดที่ใหญ่ที่สุดของหนังสือเล่มนี้ 🎓


← บทที่ 18 | สารบัญ | กลับสารบัญหลัก →