โหมดมืด
บทที่ 19 — Polyglot (พ็อล-ลี-ก็อต = "พูดหลายภาษา") Example End-to-End
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 confirmedPart 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 → ได้ topicinventory.events(รายละเอียด Outbox + Debezium ดู บทที่ 12 §CDC)สำหรับผู้ที่ต้องการ config detail: Debezium ใช้ SMT (
io.debezium.transforms.outbox.EventRouter) อ่าน columnaggregate_type = "Inventory"แล้ว route topic ตาม patternroute.topic.replacement = ${routedByValue}.events→ topicinventory.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+) — ทางเลือกอื่นคือ CustomJsonDeserializerใน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ลงในpayloadJSON ตอน 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 เดียว:
- อ่าน event จาก topic
orders.events - กัน event ซ้ำ (idempotent) ด้วยตาราง inbox
- จองของ —
UPDATE ... WHERE on_hand - reserved >= qty - เขียน 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:
@app.on_event("startup"/"shutdown")deprecated ตั้งแต่ FastAPI 0.93 (2023) → ใช้lifespanasync context manager แทน- ห้ามเรียก sync API ใน async function —
confluent_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-kafkasync ได้): ห่อด้วยasyncio.to_thread:pythonmsg = 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: 100Part 7.5: Shared Contracts — gRPC + Protobuf ข้ามภาษา
ระบบ polyglot 4 ภาษาจะอยู่ด้วยกันได้ต้องมี contract กลาง. JSON ad-hoc = ของจริงพังเร็ว — ใช้ schema ภาษากลางแล้ว generate code ทุกภาษา
ตัวเลือกมาตรฐาน 2026:
| Tool | ใช้ตอนไหน |
|---|---|
| Protocol Buffers v3 + gRPC | sync 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 Registry | event payload ใน Kafka (compact + schema evolution rules ดี) |
| JSON Schema + AsyncAPI | doc + 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 ก่อน mergebuf 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 alertPart 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 flowReliability
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 × 3Observability (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 alertsSecurity
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 bountyDeployment
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 tablePart 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
- End-to-end flow ของ order ตั้งแต่ user click จนเห็น email confirmation?
- Outbox + Inbox pattern ทำงานยังไงร่วมกันใน example นี้?
- ทำไม Order Service ใช้ Java, Inventory ใช้ Go, Notification ใช้ Node, Recommend ใช้ Python?
- Istio ทำอะไรใน example นี้?
- ถ้า Kafka down ระบบยังทำงานได้ไหม?
- ถ้า Inventory Service ตาย — order จะเกิดอะไร?
- Trace 1 request วิ่งผ่านกี่ service?
- Canary deploy ของ Order Service ทำงานยังไง?
- ทำไม Notification ใช้ทั้ง Kafka + RabbitMQ?
- ระบบนี้ 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
ขอแสดงความยินดี — คุณจบหมวดที่ใหญ่ที่สุดของหนังสือเล่มนี้ 🎓