Kinesis Data Streams สำหรับการประมวลผลเหตุการณ์แบบเรียลไทม์
สร้างและใช้สตรีมเหตุการณ์ปริมาณงานสูงด้วย Kinesis Data Streams จัดการชาร์ดเพื่อควบคุมปริมาณงาน และใช้ Lambda เป็นผู้บริโภคข้อมูล
Kinesis Data Streams สำหรับการประมวลผลเหตุการณ์แบบเรียลไทม์ เป็นบทเรียน AWS Solutions Architect ฟรีบน CoddyKit นี่คือบทเรียนที่ 3 จากทั้งหมด 4 บทเรียน คุณสามารถอ่านบทเรียนทั้งหมดด้านล่างฟรี — จากนั้นลองปฏิบัติด้วยตัวคุณเองในเบราว์เซอร์พร้อมตัวแก้ไขโค้ดในตัวและติวเตอร์ AI ตลอด 24/7 บทเรียนนี้เป็นส่วนหนึ่งของเส้นทางการเรียน AWS Solutions Architect และความก้าวหน้าของคุณจะซิงค์ข้ามเว็บและแอป CoddyKit คอร์ส AWS Solutions Architect มีบทเรียนทั้งหมด 4 บทเรียน
แนวคิดหลักของ Kinesis Data Streams
Kinesis Data Streams (KDS) คือบริการสตรีมข้อมูลแบบเรียลไทม์ที่มีความทนทานและรักษาลำดับข้อมูล ข้อมูลจะถูกจัดเก็บใน สตรีม ซึ่งประกอบด้วย ชาร์ด อย่างน้อยหนึ่งรายการ แต่ละชาร์ดคือชุดระเบียนข้อมูลที่มีลำดับ ผู้ผลิตจะส่งระเบียนไปยังชาร์ดโดยใช้ คีย์พาร์ติชัน ซึ่งกำหนดว่าระเบียนจะไปยังชาร์ดใด ผู้บริโภคจะอ่านระเบียนจากชาร์ดและประมวลผลตามลำดับที่ระเบียนมาถึงภายในชาร์ดนั้น
ความจุของชาร์ดและขีดจำกัดปริมาณงาน
แต่ละชาร์ดรองรับ ปริมาณงานการเขียน 1 MB/s หรือ 1,000 ระเบียน/s และ ปริมาณงานการอ่าน 2 MB/s (ใช้ร่วมกันระหว่างผู้บริโภคแบบมาตรฐานทั้งหมดในชาร์ดนั้น) ความจุรวมของสตรีมจะเพิ่มขึ้นตามจำนวนชาร์ดแบบเชิงเส้น ใช้สูตร: shards_needed = max(write_MB_per_s / 1, read_MB_per_s / 2) หากผู้ผลิตแตะขีดจำกัดการเขียน จะเกิดข้อผิดพลาด ProvisionedThroughputExceededException — แก้ไขโดยแบ่งชาร์ดหรือกระจายคีย์พาร์ติชันให้สม่ำเสมอยิ่งขึ้น
# Calculate shards needed for a stream:
# - Ingest rate: 5 MB/s writes
# - Read rate: 3 consumers x 2 MB/s = 6 MB/s reads
# shards = max(5/1, 6/2) = max(5, 3) = 5 shards needed
aws kinesis create-stream \
--stream-name iot-telemetry \
--shard-count 5คีย์พาร์ติชันและการกระจายข้อมูล
คีย์พาร์ติชัน คือสตริงที่ Kinesis นำไปแฮช (MD5) เพื่อกำหนดว่าระเบียนจะไปยังชาร์ดใด คีย์พาร์ติชันที่เลือกอย่างเหมาะสมจะกระจายระเบียนไปยังชาร์ดต่าง ๆ อย่างสม่ำเสมอ (การป้องกันชาร์ดร้อน) สำหรับ IoT ให้ใช้ ID อุปกรณ์ สำหรับสตรีมการคลิก ให้ใช้ ID เซสชันหรือ ID ผู้ใช้ หลีกเลี่ยงคีย์ที่มีค่าหลากหลายน้อย (เช่น ชื่อประเทศที่มีเพียง 5 ค่า) เนื่องจากจะทำให้เกิดชาร์ดร้อน ซึ่งชาร์ดหนึ่งได้รับปริมาณการเขียนมากเกินสัดส่วน ขณะที่ชาร์ดอื่นไม่มีงาน
import boto3, json
client = boto3.client('kinesis', region_name='us-east-1')
# Good: use device_id as partition key for even distribution
event = {'deviceId': 'sensor-42', 'temp': 23.5, 'ts': '2024-01-15T10:00:00Z'}
client.put_record(
StreamName='iot-telemetry',
Data=json.dumps(event),
PartitionKey='sensor-42' # high-cardinality -> even distribution
)ผู้บริโภคแบบมาตรฐานเทียบกับ Enhanced Fan-Out
ผู้บริโภคแบบมาตรฐานจะแชร์ปริมาณงานการอ่าน 2 MB/s ต่อชาร์ดโดยใช้ GetRecords และการดึงข้อมูลเป็นระยะ หากมีผู้บริโภค 3 รายในชาร์ดเดียว และแต่ละรายต้องการ 2 MB/s ผู้บริโภคเหล่านั้นจะถูกจำกัดปริมาณการใช้งาน เนื่องจากต้องแชร์ความจุรวม 2 MB/s Enhanced Fan-Out (EFO) จัดสรรท่อการอ่านเฉพาะขนาด 2 MB/s ให้ผู้บริโภคที่ลงทะเบียนแต่ละราย ผ่านการเชื่อมต่อแบบพุช HTTP/2 ที่คงอยู่ (SubscribeToShard) EFO มีค่าใช้จ่ายต่อผู้บริโภคต่อชาร์ดต่อชั่วโมง แต่กำจัดการแย่งใช้ความจุการอ่านได้ทั้งหมด
# Register an Enhanced Fan-Out consumer
aws kinesis register-stream-consumer \
--stream-arn arn:aws:kinesis:us-east-1:123456789012:stream/iot-telemetry \
--consumer-name real-time-analytics
# List registered consumers
aws kinesis list-stream-consumers \
--stream-arn arn:aws:kinesis:us-east-1:123456789012:stream/iot-telemetryLambda ในฐานะผู้บริโภค Kinesis
Lambda ผสานรวมกับ Kinesis Data Streams ได้โดยตรงผ่าน การแมปแหล่งเหตุการณ์ Lambda จะดึงข้อมูลจากสตรีม อ่านชุดระเบียน แล้วเรียกใช้ฟังก์ชันของคุณ กำหนดค่า BatchSize (1–10,000 ระเบียน), StartingPosition (TRIM_HORIZON สำหรับข้อมูลเก่าที่สุด, LATEST สำหรับข้อมูลใหม่ที่สุด) และ BisectBatchOnFunctionError เพื่อแบ่งชุดระเบียนที่ล้มเหลว ปัจจัยการทำงานแบบขนาน (1–10) ช่วยให้ Lambda เปิดใช้การเรียกใช้พร้อมกันหลายรายการต่อชาร์ด เพื่อรองรับสตรีมที่มีข้อมูลเข้ามาอย่างรวดเร็ว
# Create a Lambda event source mapping for Kinesis
aws lambda create-event-source-mapping \
--function-name ProcessIoTEvents \
--event-source-arn arn:aws:kinesis:us-east-1:123456789012:stream/iot-telemetry \
--batch-size 100 \
--starting-position LATEST \
--parallelization-factor 5 \
--bisect-batch-on-function-error true \
--destination-config '{
"OnFailure": {
"Destination": "arn:aws:sqs:us-east-1:123456789012:kinesis-dlq"
}
}'Kinesis Client Library (KCL)
Kinesis Client Library (KCL) คือเฟรมเวิร์กแอปพลิเคชันสำหรับสร้างผู้บริโภค Kinesis ที่มีความทนทานด้วย Java (หรือหลายภาษาผ่าน MultiLangDaemon) KCL จัดการการแจกแจงชาร์ด การจัดการสิทธิ์ครอบครอง (กระจายชาร์ดระหว่างอินสแตนซ์ผู้ปฏิบัติงาน) การบันทึกจุดตรวจสอบความคืบหน้าไปยัง DynamoDB และการแยกหรือรวมชาร์ดอย่างราบรื่น ผู้ปฏิบัติงาน KCL แต่ละรายจะประมวลผลชาร์ดอย่างน้อยหนึ่งรายการ และ KCL จะปรับสมดุลชาร์ดโดยอัตโนมัติเมื่อเพิ่มผู้ปฏิบัติงานหรือเมื่อผู้ปฏิบัติงานล้มเหลว โดยทั่วไป KCL เป็นตัวเลือกที่เหมาะสมกว่าการดึงข้อมูลด้วย SDK โดยตรงสำหรับแอปพลิเคชันผู้บริโภคในระบบจริง
# KCL stores checkpoints in a DynamoDB table automatically
# Each shard has one row tracking the last successfully processed sequence number
# KCL lease table structure:
# leaseKey (shardId) | checkpoint (sequenceNumber) | leaseOwner (workerId)
#
# To start a KCL application (pseudocode):
# KinesisClientLibConfiguration config = new KinesisClientLibConfiguration(
# 'iot-app', 'iot-telemetry', credentialsProvider, 'worker-1');
# Worker worker = new Worker.Builder().config(config).recordProcessorFactory(factory).build();
# worker.run();การเก็บรักษาข้อมูลและการเล่นข้อมูลซ้ำ
Kinesis Data Streams จัดเก็บระเบียนเป็นเวลา 24 ชั่วโมงตามค่าเริ่มต้น (ขยายได้ถึง 7 วัน หรือสูงสุด 365 วันด้วยการเก็บรักษาระยะยาว โดยมีค่าใช้จ่ายเพิ่มเติม) ต่างจาก SQS ระเบียนที่ถูกอ่านแล้วจะ ไม่ถูกลบ หลังการใช้งาน แต่จะยังคงพร้อมใช้งานจนกว่าระยะเวลาเก็บรักษาจะหมดลง คุณสมบัตินี้ช่วยให้ผู้บริโภคหลายรายอ่านระเบียนเดียวกันได้อย่างอิสระ และทำให้สามารถ เล่นข้อมูลซ้ำ โดยรีเซ็ตจุดตรวจสอบของผู้บริโภคไปยัง sequence number ก่อนหน้า ซึ่งมีประโยชน์อย่างยิ่งสำหรับการแก้ไขข้อผิดพลาดหรือเติมข้อมูลให้บริการใหม่
# Extend stream retention to 7 days
aws kinesis increase-stream-retention-period \
--stream-name iot-telemetry \
--retention-period-hours 168
# Get records from the oldest available record (replay)
SHARD_ITERATOR=$(aws kinesis get-shard-iterator \
--stream-name iot-telemetry \
--shard-id shardId-000000000000 \
--shard-iterator-type TRIM_HORIZON \
--query 'ShardIterator' --output text)
aws kinesis get-records --shard-iterator $SHARD_ITERATOR --limit 100โหมดความจุแบบ On-Demand
Kinesis Data Streams รองรับโหมดความจุสองแบบ โหมด Provisioned: คุณจัดการจำนวนชาร์ดด้วยตนเองและชำระค่าบริการตามชาร์ดต่อชั่วโมง โหมด On-Demand: Kinesis ปรับขนาดความจุของชาร์ดโดยอัตโนมัติตามปริมาณงานขาเข้า (ตามค่าเริ่มต้น สูงสุด 200 MB/s สำหรับการเขียน และ 400 MB/s สำหรับการอ่าน) และคุณชำระค่าบริการตาม GB ของข้อมูลที่เขียนและดึงข้อมูล โหมด On-Demand เหมาะสำหรับรูปแบบการรับส่งข้อมูลที่เปลี่ยนแปลงหรือคาดเดาไม่ได้ ซึ่งคุณไม่ต้องการจัดการการปรับขนาดชาร์ด
# Switch an existing stream to On-Demand mode
aws kinesis update-stream-mode \
--stream-arn arn:aws:kinesis:us-east-1:123456789012:stream/iot-telemetry \
--stream-mode-details StreamMode=ON_DEMANDการรับประกันลำดับข้อมูลภายในชาร์ด
Kinesis รับประกัน ลำดับข้อมูลภายในชาร์ด — ระเบียนที่มีคีย์พาร์ติชันเดียวกันจะไปยังชาร์ดเดียวกันเสมอ และจะถูกอ่านตามลำดับที่เขียน อย่างไรก็ตาม ไม่มีการรับประกันลำดับข้อมูลระหว่างชาร์ด หากแอปพลิเคชันต้องการลำดับข้อมูลรวมของระเบียนทั้งหมด ให้ใช้ชาร์ดเดียว (ซึ่งจำกัดปริมาณงานไว้ที่ 1 MB/s) หรือออกแบบใหม่ให้ต้องการลำดับเฉพาะภายในกลุ่มคีย์พาร์ติชันเท่านั้น (เช่น ลำดับข้อมูลแยกตามอุปกรณ์) ประเด็นนี้เป็นข้อแตกต่างที่มักออกสอบเมื่อเปรียบเทียบกับ SQS FIFO ซึ่งมีการกำจัดข้อมูลซ้ำและการจัดลำดับอย่างเคร่งครัด
ความปลอดภัย: การเข้ารหัสและ VPC
Kinesis Data Streams เข้ารหัสระเบียน ฝั่งเซิร์ฟเวอร์ขณะจัดเก็บ โดยใช้ AWS KMS (CMK หรือคีย์ที่ AWS จัดการ) เมื่อเปิดใช้การเข้ารหัสฝั่งเซิร์ฟเวอร์ ข้อมูลทั้งหมดระหว่างการส่งจะถูกเข้ารหัสด้วย TLS สำหรับแอปพลิเคชันที่ทำงานใน VPC และไม่ควรส่งข้อมูลสตรีมผ่านอินเทอร์เน็ตสาธารณะ ให้ใช้ VPC Interface Endpoint (PrivateLink) สำหรับ Kinesis เพื่อให้การรับส่งข้อมูลอยู่ภายในโครงข่ายหลักของ AWS ทั้งหมด ซึ่งสำคัญสำหรับงานที่อยู่ภายใต้ข้อกำกับดูแล
# Enable server-side encryption on a Kinesis stream
aws kinesis start-stream-encryption \
--stream-name iot-telemetry \
--encryption-type KMS \
--key-id arn:aws:kms:us-east-1:123456789012:key/mrk-abc123Kinesis Data Streams เทียบกับ SQS และ Kafka
สำหรับการสอบ SAA-C03 ให้เปรียบเทียบ Kinesis Data Streams กับทางเลือกอื่น Kinesis เทียบกับ SQS: Kinesis รักษาลำดับข้อมูลภายในชาร์ดและรองรับผู้บริโภคหลายรายที่อ่านข้อมูลเดียวกัน ส่วน SQS จะลบข้อความหลังการใช้งาน Kinesis เทียบกับ MSK (Kafka): MSK คือ Apache Kafka ที่มีการจัดการ — ให้ใช้เมื่อจำเป็นต้องรองรับความเข้ากันได้กับโพรโทคอล Kafka การกำหนดค่าท็อปปิกขั้นสูง หรือการย้ายระบบจาก Kafka ภายในองค์กร ให้ใช้ Kinesis สำหรับการสตรีมที่ใช้บริการของ AWS โดยเฉพาะและผสานรวมกับ Lambda, Firehose และ Flink ได้แน่นแฟ้นกว่า เลือก Kinesis เว้นแต่โจทย์จะระบุ Kafka หรือข้อกำหนดด้านความเข้ากันได้กับ Kafka โดยเฉพาะ
ตรวจสอบความเข้าใจอย่างรวดเร็ว
ทดสอบความเข้าใจแนวคิด AWS Solutions Architect (SAA-C03) จากบทเรียนนี้
สรุปบทเรียน
ในบทเรียนนี้ คุณได้เรียนรู้ว่า Kinesis Data Streams ให้บริการสตรีมที่แบ่งเป็นชาร์ด มีลำดับ มีความทนทาน เก็บรักษาข้อมูลตามค่าเริ่มต้น 24 ชั่วโมง และรองรับการเล่นข้อมูลซ้ำ Enhanced Fan-Out จัดสรรความเร็ว 2 MB/s ต่อชาร์ดให้ผู้บริโภคแต่ละรายโดยเฉพาะ จึงกำจัดการแย่งใช้ความจุการอ่าน และ โหมด On-Demand ปรับขนาดชาร์ดโดยอัตโนมัติเพื่อรองรับการรับส่งข้อมูลที่คาดเดาไม่ได้ ต่อไป เราจะศึกษาแพตเทิร์น choreography เทียบกับ orchestration ในสถาปัตยกรรมที่ขับเคลื่อนด้วยเหตุการณ์
คำถามที่พบบ่อย
บทเรียน “Kinesis Data Streams สำหรับการประมวลผลเหตุการณ์แบบเรียลไทม์” ฟรีหรือไม่
ใช่ — ข้อความเต็มของ “Kinesis Data Streams สำหรับการประมวลผลเหตุการณ์แบบเรียลไทม์” ฟรีให้อ่านที่นี่บนเว็บ เพื่อปฏิบัติแบบโต้ตอบ (ตัวแก้ไขโค้ดในตัวและติวเตอร์ AI ตลอด 24/7) และปลดล็อคส่วนที่เหลือของคอร์ส AWS Solutions Architect ให้อัปเกรดเป็น CoddyKit PRO คอร์ส AWS Solutions Architect มีบทเรียนทั้งหมด 4 บทเรียน
คุณจะเรียนรู้อะไรในบทเรียน “Kinesis Data Streams สำหรับการประมวลผลเหตุการณ์แบบเรียลไทม์”
สร้างและใช้สตรีมเหตุการณ์ปริมาณงานสูงด้วย Kinesis Data Streams จัดการชาร์ดเพื่อควบคุมปริมาณงาน และใช้ Lambda เป็นผู้บริโภคข้อมูล คุณปฏิบัติ AWS Solutions Architect ด้วยโค้ดที่ใช้งานได้จริงที่คุณเรียกใช้โดยตรงในเบราว์เซอร์ และติวเตอร์ AI ตลอด 24/7 ตอบคำถามของคุณขณะที่คุณไปผ่านบทเรียน
คุณต้องมีประสบการณ์ก่อนที่จะเริ่มเรียน AWS Solutions Architect หรือไม่
ไม่จำเป็นต้องมีประสบการณ์มาก่อน AWS Solutions Architect บน CoddyKit ออกแบบมาสำหรับผู้เริ่มต้นไปจนถึงผู้เรียนขั้นสูง คุณสามารถเริ่มต้นที่นี่หรือเริ่มจากตัวแรกและเรียนด้วยความเร็วของคุณเอง นี่คือบทเรียนที่ 3 จากทั้งหมด 4 บทเรียน
บทเรียน “Kinesis Data Streams สำหรับการประมวลผลเหตุการณ์แบบเรียลไทม์” ใช้เวลานานแค่ไหน
บทเรียน CoddyKit ส่วนใหญ่ใช้เวลาประมาณ 5–10 นาที แต่ละบทเรียนจึงสั้นและเป็นแบบโต้ตอบ คุณสามารถก้าวหน้าอย่างต่อเนื่องและกลับมาเรียนต่อจากตรงที่เพิ่งหยุดบนเว็บและแอปได้เลย
ฉันเขียนและรันโค้ดในบทเรียน AWS Solutions Architect นี้ได้ไหม
ได้ บทเรียน AWS Solutions Architect ทุกบทมีตัวแก้ไขโค้ดในตัว คุณจึงเขียนและรันโค้ดจริงได้เลยในเบราว์เซอร์ และได้รับข้อเสนอแนะจาก AI ในทันที — ไม่ต้องติดตั้งในเครื่องของคุณ
บทเรียนทั้งหมดในหลักสูตรนี้
- EventBridge: บัสเหตุการณ์และกฎ
- Step Functions: การประสานกระบวนงานแบบไร้เซิร์ฟเวอร์
- Kinesis Data Streams สำหรับการประมวลผลเหตุการณ์แบบเรียลไทม์
- รูปแบบการประสานงานกับการควบคุมจากศูนย์กลาง