Kinesis Streams, Firehose และการวิเคราะห์แบบเรียลไทม์
รับข้อมูลสตรีมด้วย Kinesis Data Streams ส่งข้อมูลไปยัง S3 หรือ Redshift ด้วย Firehose และวิเคราะห์แบบเรียลไทม์ด้วย Managed Service for Apache Flink
Kinesis Streams, Firehose และการวิเคราะห์แบบเรียลไทม์ เป็นบทเรียน AWS Solutions Architect ฟรีบน CoddyKit นี่คือบทเรียนที่ 4 จากทั้งหมด 4 บทเรียน คุณสามารถอ่านบทเรียนทั้งหมดด้านล่างฟรี — จากนั้นลองปฏิบัติด้วยตัวคุณเองในเบราว์เซอร์พร้อมตัวแก้ไขโค้ดในตัวและติวเตอร์ AI ตลอด 24/7 บทเรียนนี้เป็นส่วนหนึ่งของเส้นทางการเรียน AWS Solutions Architect และความก้าวหน้าของคุณจะซิงค์ข้ามเว็บและแอป CoddyKit คอร์ส AWS Solutions Architect มีบทเรียนทั้งหมด 4 บทเรียน
ภาพรวมตระกูล Kinesis
Amazon Kinesis คือตระกูลบริการสำหรับรวบรวม ประมวลผล และวิเคราะห์ข้อมูลสตรีมมิงแบบเรียลไทม์ บริการหลักสามรายการ ได้แก่ Kinesis Data Streams (การประมวลผลแบบกำหนดเองที่มีเวลาแฝงต่ำ), Kinesis Data Firehose (การส่งข้อมูลที่มีการจัดการเต็มรูปแบบไปยัง S3/Redshift/OpenSearch) และ Managed Service for Apache Flink (เดิมคือ Kinesis Data Analytics) สำหรับ SQL และการประมวลผลด้วย Flink แบบเรียลไทม์ แต่ละบริการเหมาะกับจุดประสงค์ที่แตกต่างกันในกระบวนการสตรีมข้อมูล
สถาปัตยกรรม Kinesis Data Streams
Kinesis Data Stream คือบันทึกข้อมูลแบบถาวรและมีลำดับ ซึ่งแบ่งออกเป็น ชาร์ด แต่ละชาร์ดรองรับปริมาณการเขียน 1 MB/s และการอ่าน 2 MB/s ข้อมูลจะถูกเก็บรักษาไว้ 24 ชั่วโมงเป็นค่าเริ่มต้น (ขยายได้ถึง 7 วันหรือ 365 วัน) ผู้ผลิตเขียนระเบียนลงในสตรีม ส่วนผู้ใช้ — Lambda, แอปพลิเคชัน KCL, Firehose หรือ Flink — อ่านข้อมูลจากชาร์ดอย่างน้อยหนึ่งชาร์ดแบบขนาน ระเบียนจะไม่สามารถแก้ไขได้หลังจากเขียนแล้ว
# Create a Kinesis Data Stream with 4 shards
aws kinesis create-stream \
--stream-name clickstream \
--shard-count 4
# Put a record into the stream
aws kinesis put-record \
--stream-name clickstream \
--partition-key 'user-123' \
--data 'eyJldmVudCI6ICJjbGljayJ9'ชาร์ด ปริมาณงาน และการปรับขนาด
จำนวนชาร์ดเป็นตัวกำหนดปริมาณงานรวมของสตรีม คุณสามารถ แยกชาร์ด เพื่อเพิ่มปริมาณงานเป็นสองเท่า หรือ รวมชาร์ดสองชาร์ด เพื่อลดค่าใช้จ่าย ใช้ Enhanced Fan-Out เพื่อให้ผู้ใช้ที่ลงทะเบียนแต่ละรายมีปริมาณงานการอ่านของตนเองที่ 2 MB/s โดยไม่ขึ้นต่อผู้ใช้รายอื่น ช่วยขจัดการจำกัดการอ่านเมื่อมีหลายแอปพลิเคชันใช้สตรีมเดียวกัน ตรวจสอบ GetRecords.IteratorAgeMilliseconds เพื่อตรวจจับความล่าช้าของผู้ใช้
# Split shard to increase throughput
aws kinesis split-shard \
--stream-name clickstream \
--shard-to-split shardId-000000000001 \
--new-starting-hash-key 170141183460469231731687303715884105728
# Register an enhanced fan-out consumer
aws kinesis register-stream-consumer \
--stream-arn arn:aws:kinesis:us-east-1:123456789012:stream/clickstream \
--consumer-name analytics-appKinesis Data Firehose: การส่งข้อมูลที่มีการจัดการ
Kinesis Data Firehose คือบริการที่มีการจัดการเต็มรูปแบบสำหรับรวบรวม แปลง และส่งข้อมูลสตรีมมิงไปยังปลายทางต่าง ๆ รวมถึง S3, Amazon Redshift, Amazon OpenSearch Service, Splunk และปลายทาง HTTP ไม่ต้องจัดการชาร์ด เพราะ Firehose จะปรับขนาดโดยอัตโนมัติ คุณกำหนดค่า ขนาดบัฟเฟอร์ (1–128 MB) และ ช่วงเวลาบัฟเฟอร์ (60–900 วินาที) โดย Firehose จะส่งข้อมูลเมื่อถึงขีดจำกัดใดขีดจำกัดหนึ่งก่อน
# Create a Firehose delivery stream to S3
aws firehose create-delivery-stream \
--delivery-stream-name clickstream-to-s3 \
--s3-destination-configuration '{
"RoleARN": "arn:aws:iam::123456789012:role/FirehoseRole",
"BucketARN": "arn:aws:s3:::my-data-lake-123",
"Prefix": "landing/clickstream/year=!{timestamp:yyyy}/month=!{timestamp:MM}/",
"BufferingHints": {"SizeInMBs": 64, "IntervalInSeconds": 300},
"CompressionFormat": "GZIP"
}'การแปลงข้อมูล Firehose ด้วย Lambda
Firehose สามารถเรียกใช้ ฟังก์ชัน Lambda กับระเบียนแต่ละชุดก่อนส่งข้อมูล เพื่อแปลง เพิ่มข้อมูล หรือกรองข้อมูลระหว่างทาง กรณีใช้งานทั่วไป ได้แก่ การแปลง JSON เป็น Parquet (ผ่านสคีมา Glue) การปกปิดฟิลด์ PII หรือการทิ้งเหตุการณ์ที่มีคุณค่าต่ำ ระเบียนที่แปลงไม่สำเร็จสามารถส่งไปยังคำนำหน้าเอาต์พุตข้อผิดพลาด S3 แยกต่างหากเพื่อประมวลผลใหม่ได้ตามต้องการ จึงไม่มีข้อมูลสูญหาย
# Lambda transform function signature for Firehose
def lambda_handler(event, context):
output = []
for record in event['records']:
import base64, json
payload = json.loads(base64.b64decode(record['data']))
# Drop events with no user_id
if not payload.get('user_id'):
output.append({'recordId': record['recordId'], 'result': 'Dropped', 'data': record['data']})
else:
output.append({'recordId': record['recordId'], 'result': 'Ok', 'data': record['data']})
return {'records': output}บริการที่มีการจัดการสำหรับ Apache Flink
Amazon Managed Service for Apache Flink (เดิมคือ Kinesis Data Analytics) เรียกใช้แอปพลิเคชัน Apache Flink บนโครงสร้างพื้นฐานที่มีการจัดการเต็มรูปแบบ ใช้บริการนี้สำหรับการวิเคราะห์แบบมีสถานะและเรียลไทม์ เช่น การรวมค่าจากหน้าต่างเลื่อน การตรวจจับความผิดปกติ การจับคู่รูปแบบในลำดับเหตุการณ์ และการ JOIN ข้อมูลสตรีมมิงกับตารางอ้างอิง คุณเขียนโค้ด Flink ด้วย Java, Python หรือ Scala ส่วน Flink จะจัดการจุดตรวจสอบและสถานะแบบประมวลผลครั้งเดียวอย่างแน่นอน
# Flink SQL-style tumbling window (conceptual)
# Count page views per URL every 5 minutes
CREATE TABLE clickstream (
url STRING,
event_time TIMESTAMP(3),
WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND
) WITH ('connector' = 'kinesis', 'stream' = 'clickstream', ...);
SELECT
url,
COUNT(*) AS views,
TUMBLE_START(event_time, INTERVAL '5' MINUTE) AS window_start
FROM clickstream
GROUP BY url, TUMBLE(event_time, INTERVAL '5' MINUTE);การเลือกระหว่าง Streams กับ Firehose
สำหรับข้อสอบ SAA-C03 คุณควรทราบเกณฑ์การตัดสินใจ ให้ใช้ Kinesis Data Streams เมื่อต้องการเวลาแฝงต่ำกว่าหนึ่งวินาที มีผู้ใช้หลายรายอ่านข้อมูลพร้อมกัน หรือต้องการตรรกะการประมวลผลแบบกำหนดเองพร้อมควบคุมการเก็บรักษาและการเล่นข้อมูลซ้ำได้อย่างเต็มที่ ให้ใช้ Firehose เมื่อต้องการเพียงส่งข้อมูลสตรีมมิงไปยัง S3, Redshift หรือ OpenSearch อย่างน่าเชื่อถือ โดยใช้โค้ดเพียงเล็กน้อย มีการแปลงข้อมูลระหว่างทางตามต้องการ และปรับขนาดอัตโนมัติ แม้จะมีเวลาแฝงสูงกว่า (อย่างน้อย 60 วินาที)
Kinesis เทียบกับ SQS: ตัวเลือกคลาสสิกในข้อสอบ
คำถามทั่วไปในข้อสอบ SAA-C03 คือให้เลือกใช้ Kinesis หรือ SQS ความแตกต่างสำคัญมีดังนี้: Kinesis รักษา ลำดับข้อความภายในชาร์ด รองรับผู้ใช้หลายรายที่อ่านข้อมูลเดียวกันพร้อมกัน และเก็บรักษาระเบียนไว้เพื่อเล่นข้อมูลซ้ำได้ ส่วน SQS จะลบข้อความหลังจากถูกใช้แล้ว (ไม่สามารถเล่นซ้ำได้) โหมด FIFO รับประกันลำดับที่เคร่งครัด และเหมาะกว่าสำหรับแยกการทำงานของไมโครเซอร์วิส หากสถานการณ์กล่าวถึงการวิเคราะห์แบบเรียลไทม์หรือการเล่นข้อมูลซ้ำ ให้เลือก Kinesis
ผู้ผลิต Kinesis: SDK และ KPL
Kinesis Producer Library (KPL) คือไคลเอ็นต์ประสิทธิภาพสูงสำหรับเขียนข้อมูลไปยัง Kinesis Data Streams จากแอปพลิเคชัน KPL จะ รวม ระเบียนขนาดเล็กหลายรายการเป็นการเรียก API เดียวโดยอัตโนมัติ (สูงสุด 1 MB) และจัดการการลองใหม่พร้อมการหน่วงเวลาแบบเพิ่มขึ้น วิธีนี้ช่วยลดค่าใช้จ่าย PUT ต่อระเบียนและเพิ่มปริมาณงานต่อชาร์ดได้อย่างมาก ใช้ KPL กับผู้ผลิตข้อมูลปริมาณสูง เช่น สตรีมการคลิกบนเว็บไซต์ เซนเซอร์ IoT หรือกระบวนการส่งบันทึก
# Basic KPL usage (Java pseudocode, no backticks)
KinesisProducer producer = new KinesisProducer();
byte[] data = 'hello world'.getBytes();
ByteBuffer buf = ByteBuffer.wrap(data);
// addUserRecord handles aggregation and retry internally
ListenableFuture future = producer.addUserRecord('clickstream', 'partitionKey', buf);
Futures.addCallback(future, new FutureCallback() { ... });รูปแบบการส่งข้อมูลจาก Firehose ไปยัง Redshift
สถาปัตยกรรมที่พบได้ทั่วไปคือใช้ Firehose เป็นกระบวนการส่งข้อมูลที่มีการจัดการจากสตรีม Kinesis หรือผู้ผลิตข้อมูลโดยตรงไปยัง Amazon Redshift เพื่อสร้างคลังข้อมูล Firehose จะเขียนข้อมูลลงในบักเก็ต S3 ชั่วคราวก่อน จากนั้นจึงออกคำสั่ง COPY เพื่อโหลดข้อมูลเข้า Redshift วิธีนี้มีประสิทธิภาพสูงสุดสำหรับการโหลดข้อมูลสตรีมมิงจำนวนมากเข้า Redshift การเพิ่มข้อมูลทีละแถวโดยตรงใน Redshift จะช้ามากเนื่องจากมีค่าใช้จ่ายในการประมวลผลต่อแถว
# Firehose Redshift destination (CLI snippet)
--redshift-destination-configuration '{
"RoleARN": "arn:aws:iam::123456789012:role/FirehoseRole",
"ClusterJDBCURL": "jdbc:redshift://cluster.xyz.us-east-1.redshift.amazonaws.com:5439/sales",
"CopyCommand": {
"DataTableName": "clickevents",
"CopyOptions": "JSON 'auto'"
},
"Username": "firehose_user",
"Password": "{{resolve:secretsmanager:redshift-pw}}",
"S3Configuration": {
"RoleARN": "...",
"BucketARN": "arn:aws:s3:::firehose-staging"
}
}'การตรวจสอบ Kinesis ด้วย CloudWatch
เมตริก CloudWatch ที่สำคัญสำหรับ Kinesis ได้แก่ IncomingBytes และ IncomingRecords สำหรับวัดปริมาณงานของผู้ผลิต GetRecords.IteratorAgeMilliseconds สำหรับวัดความล่าช้าของผู้ใช้ (ค่าที่สูงหมายความว่าผู้ใช้ประมวลผลข้อมูลไม่ทัน) และ WriteProvisionedThroughputExceeded สำหรับตรวจจับว่าผู้ผลิตกำลังชนขีดจำกัดของชาร์ดหรือไม่ ตั้งการแจ้งเตือน CloudWatch สำหรับอายุของตัววนซ้ำและปริมาณงานที่จัดสรรเกิน เพื่อทริกเกอร์การปรับขนาดชาร์ดอัตโนมัติผ่าน Application Auto Scaling
# CloudWatch alarm on high consumer lag
aws cloudwatch put-metric-alarm \
--alarm-name kinesis-high-lag \
--metric-name GetRecords.IteratorAgeMilliseconds \
--namespace AWS/Kinesis \
--dimensions Name=StreamName,Value=clickstream \
--statistic Maximum \
--period 60 \
--threshold 60000 \
--comparison-operator GreaterThanThreshold \
--evaluation-periods 3 \
--alarm-actions arn:aws:sns:us-east-1:123456789012:ops-alertsตรวจสอบความเข้าใจ
ทดสอบความเข้าใจแนวคิด AWS Solutions Architect (SAA-C03) จากบทเรียนนี้
สรุปบทเรียน
ในบทเรียนนี้ คุณได้เรียนรู้ว่า Kinesis Data Streams ให้บริการสตรีมข้อมูลแบบแบ่งชาร์ดที่มีลำดับและความทนทาน รองรับผู้ใช้หลายรายและการเล่นข้อมูลซ้ำ, Kinesis Data Firehose ให้บริการส่งข้อมูลไปยัง S3, Redshift และ OpenSearch แบบมีการจัดการเต็มรูปแบบโดยไม่ต้องเขียนโค้ด และ Managed Service for Apache Flink ช่วยให้วิเคราะห์ข้อมูลสตรีมแบบเรียลไทม์ที่มีสถานะได้ บทถัดไป เราจะสำรวจกลยุทธ์การย้ายระบบขึ้นคลาวด์ 7 Rs
เรียนรู้ AWS Solutions Architect ด้วย AI tutor — ฟรี
เขียนและเรียกใช้โค้ดจริงในเบราว์เซอร์ของคุณ รับความช่วยเหลือทันทีจาก AI tutor 24/7 และเรียนรู้ต่อจากที่คุณหยุดบนเว็บหรือในแอป
- คอร์ส
- 30
- บทเรียน
- 120
คำถามที่พบบ่อย
บทเรียน “Kinesis Streams, Firehose และการวิเคราะห์แบบเรียลไทม์” ฟรีหรือไม่
ใช่ — ข้อความเต็มของ “Kinesis Streams, Firehose และการวิเคราะห์แบบเรียลไทม์” ฟรีให้อ่านที่นี่บนเว็บ เพื่อปฏิบัติแบบโต้ตอบ (ตัวแก้ไขโค้ดในตัวและติวเตอร์ AI ตลอด 24/7) และปลดล็อคส่วนที่เหลือของคอร์ส AWS Solutions Architect ให้อัปเกรดเป็น CoddyKit PRO คอร์ส AWS Solutions Architect มีบทเรียนทั้งหมด 4 บทเรียน
คุณจะเรียนรู้อะไรในบทเรียน “Kinesis Streams, Firehose และการวิเคราะห์แบบเรียลไทม์”
รับข้อมูลสตรีมด้วย Kinesis Data Streams ส่งข้อมูลไปยัง S3 หรือ Redshift ด้วย Firehose และวิเคราะห์แบบเรียลไทม์ด้วย Managed Service for Apache Flink คุณปฏิบัติ AWS Solutions Architect ด้วยโค้ดที่ใช้งานได้จริงที่คุณเรียกใช้โดยตรงในเบราว์เซอร์ และติวเตอร์ AI ตลอด 24/7 ตอบคำถามของคุณขณะที่คุณไปผ่านบทเรียน
คุณต้องมีประสบการณ์ก่อนที่จะเริ่มเรียน AWS Solutions Architect หรือไม่
ไม่จำเป็นต้องมีประสบการณ์มาก่อน AWS Solutions Architect บน CoddyKit ออกแบบมาสำหรับผู้เริ่มต้นไปจนถึงผู้เรียนขั้นสูง คุณสามารถเริ่มต้นที่นี่หรือเริ่มจากตัวแรกและเรียนด้วยความเร็วของคุณเอง นี่คือบทเรียนที่ 4 จากทั้งหมด 4 บทเรียน
บทเรียน “Kinesis Streams, Firehose และการวิเคราะห์แบบเรียลไทม์” ใช้เวลานานแค่ไหน
บทเรียน CoddyKit ส่วนใหญ่ใช้เวลาประมาณ 5–10 นาที แต่ละบทเรียนจึงสั้นและเป็นแบบโต้ตอบ คุณสามารถก้าวหน้าอย่างต่อเนื่องและกลับมาเรียนต่อจากตรงที่เพิ่งหยุดบนเว็บและแอปได้เลย
ฉันเขียนและรันโค้ดในบทเรียน AWS Solutions Architect นี้ได้ไหม
ได้ บทเรียน AWS Solutions Architect ทุกบทมีตัวแก้ไขโค้ดในตัว คุณจึงเขียนและรันโค้ดจริงได้เลยในเบราว์เซอร์ และได้รับข้อเสนอแนะจาก AI ในทันที — ไม่ต้องติดตั้งในเครื่องของคุณ
บทเรียนทั้งหมดในหลักสูตรนี้
- การสร้างดาต้าเลคบน S3
- AWS Glue: ETL และแค็ตตาล็อกข้อมูล
- Amazon Athena: SQL แบบไร้เซิร์ฟเวอร์บน S3
- Kinesis Streams, Firehose และการวิเคราะห์แบบเรียลไทม์