Event source mappings และการประมวลผล stream

Serverless Applications with AWS Lambda

Claudio Canales

Senior DevOps Engineer

ตำแหน่งของ event source mappings

  • คิวและ stream มักถูก poll เป็น batch
  • Event source mappings จัดการพฤติกรรมการ poll และการรวม batch

ตำแหน่งของ event source mapping

Serverless Applications with AWS Lambda

Event source mapping คืออะไร?

  • ทรัพยากรที่เชื่อมต่อ source กับฟังก์ชันของคุณ
  • Lambda poll source เพื่อดึง record
  • Lambda เรียก handler พร้อม batch

ขั้นตอนการทำงานของ event source mapping

Serverless Applications with AWS Lambda

Push vs poll (ความสำคัญ)

Push model

  • Push source ส่ง event ทันที
  • ตัวอย่าง: การอัปโหลดไปยัง Amazon S3

Poll model

  • Lambda ตรวจสอบ poll source เอง
  • Lambda รวบรวม batch แล้วเรียก handler
  • ส่งผลต่อ latency และพฤติกรรมการลองใหม่

Push vs poll models

Serverless Applications with AWS Lambda

Queue vs stream

Queue (Amazon SQS)

  • คิวข้อความแบบ managed คล้ายกล่องจดหมาย

Stream (DynamoDB Streams)

  • บันทึกการเปลี่ยนแปลงของตาราง DynamoDB
  • ทั้งคู่ส่งข้อมูลในรูปแบบ Records แต่ความหมายต่างกัน

เปรียบเทียบ queue กับ stream

Serverless Applications with AWS Lambda

รูปแบบ batch event (แบบย่อ)

{
  "Records": [{
    "messageId": "abc-123",
    "body": "{\"order_id\": \"A-42\"}"
  }]
}
  • event ของ mapping ส่วนใหญ่เริ่มต้นด้วย Records
  • แต่ละ record มี messageId และ string body
  • แปลง body เป็น payload ของคุณเอง
Serverless Applications with AWS Lambda

Walkthrough: ประมวลผลแต่ละ record

def lambda_handler(event, context):
    records = event.get("Records", [])
    for record in records:
        body = record.get("body", "")
        print("BODY:", body)
    return {"statusCode": 200}
  • อ่าน Records โดยกำหนดค่าเริ่มต้นเป็น list ว่าง
  • วนลูปแต่ละ record และอ่าน body อย่างปลอดภัย
  • บันทึกข้อมูลที่ต้องการ แล้ว return
Serverless Applications with AWS Lambda

Walkthrough: แปลง JSON body

import json

def lambda_handler(event, context):
    record = event.get("Records", [])[0]
    payload = json.loads(record.get("body", "{}"))
    order_id = payload.get("order_id")
    print("ORDER_ID:", order_id)
    return {"statusCode": 200}
  • import json และอ่าน body พร้อมค่า default ที่ปลอดภัย
  • แปลงด้วย json.loads เพื่อให้ได้ dict
  • ดึงฟิลด์ที่ต้องการ โดยตรวจสอบก่อนเสมอ
Serverless Applications with AWS Lambda

Batching: batch size และ batch window

  • Batch size: จำนวน record ต่อการเรียกใช้งานหนึ่งครั้ง
  • Batch window: เวลาที่ Lambda รอเพื่อรวม batch
  • batch ขนาดใหญ่เพิ่ม throughput แต่ใช้งานมากขึ้นต่อครั้ง
  • สำหรับ SQS ให้ติดตาม ApproximateAgeOfOldestMessage เพื่อตรวจสอบ backlog

แผนภาพ trade-off ของ batch size

Serverless Applications with AWS Lambda

Trade-off ของ batch size

  • batch ขนาดใหญ่มีประสิทธิภาพสูง แต่อาจเพิ่ม latency
  • batch ขนาดเล็กลด latency แต่เพิ่มจำนวนการเรียกใช้งาน
  • เลือกให้เหมาะกับ workload ของคุณ

กราฟ batch size vs latency

Serverless Applications with AWS Lambda

การ scale ด้วย concurrency

  • หลาย batch ถูกประมวลผลพร้อมกัน
  • concurrency ที่สูงขึ้นเพิ่ม throughput
  • แต่ก็เพิ่มภาระให้ระบบปลายทางด้วย

การประมวลผล batch แบบ concurrent

Serverless Applications with AWS Lambda

ปกป้องระบบปลายทาง

  • concurrency มากเกินไปอาจทำให้ database หรือ API รับไม่ไหว
  • ใช้ limit เพื่อควบคุมภาระ
  • ควรแบ่งงานเป็นหน่วยเล็กและปลอดภัย

การป้องกันภาระของระบบปลายทาง

Serverless Applications with AWS Lambda

Partial batch failures

  • batch อาจมีทั้ง record ที่ดีและที่มีปัญหา
  • หากไม่รองรับ partial failure record ที่มีปัญหาจะทำให้ทั้ง batch ถูกลองใหม่
  • การจัดการ partial failure ช่วยให้ลองใหม่เฉพาะรายการที่ล้มเหลว

ขั้นตอน partial batch failure

Serverless Applications with AWS Lambda

Walkthrough: รายงาน partial failures (SQS)

def lambda_handler(event, context):
    failures = []
    for record in event.get("Records", []):
        try:
            process(record)
        except Exception:
            failures.append({"itemIdentifier": record["messageId"]})
    return {"batchItemFailures": failures}
  • รวบรวม messageId ของ record ที่ล้มเหลวใน batchItemFailures
  • Lambda จะลองใหม่เฉพาะรายการเหล่านั้น
  • record ที่ประมวลผลสำเร็จจะไม่ถูกส่งซ้ำ
Serverless Applications with AWS Lambda

สรุปสิ่งที่ได้เรียนรู้

  • Event source mappings poll คิวและ stream
  • ส่ง Records เป็น batch ไปยัง handler
  • batch size และ window แลกกันระหว่าง latency กับ throughput
  • Partial failures ป้องกันการประมวลผลซ้ำในส่วนที่สำเร็จแล้ว

แผนภาพสรุป mapping

Serverless Applications with AWS Lambda

มาฝึกกันเถอะ!

Serverless Applications with AWS Lambda

Preparing Video For Download...