Kinesis Streams, Firehose และการวิเคราะห์แบบเรียลไทม์
รับข้อมูลสตรีมด้วย Kinesis Data Streams ส่งข้อมูลไปยัง S3 หรือ Redshift ด้วย Firehose และวิเคราะห์แบบเรียลไทม์ด้วย Managed Service for Apache Flink
Kinesis Streams, Firehose และการวิเคราะห์แบบเรียลไทม์ เป็นบทเรียน Cloud & IT Cert Prep ฟรีบน CoddyKit นี่คือบทเรียนที่ 4 จากทั้งหมด 4 บทเรียน คุณสามารถอ่านบทเรียนทั้งหมดด้านล่างฟรี — จากนั้นลองปฏิบัติด้วยตัวคุณเองในเบราว์เซอร์พร้อมตัวแก้ไขโค้ดในตัวและติวเตอร์ AI ตลอด 24/7 บทเรียนนี้เป็นส่วนหนึ่งของเส้นทางการเรียน Cloud & IT Cert Prep และความก้าวหน้าของคุณจะซิงค์ข้ามเว็บและแอป CoddyKit คอร์ส Cloud & IT Cert Prep มีบทเรียนทั้งหมด 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
คำถามที่พบบ่อย
บทเรียน “Kinesis Streams, Firehose และการวิเคราะห์แบบเรียลไทม์” ฟรีหรือไม่
ใช่ — ข้อความเต็มของ “Kinesis Streams, Firehose และการวิเคราะห์แบบเรียลไทม์” ฟรีให้อ่านที่นี่บนเว็บ เพื่อปฏิบัติแบบโต้ตอบ (ตัวแก้ไขโค้ดในตัวและติวเตอร์ AI ตลอด 24/7) และปลดล็อคส่วนที่เหลือของคอร์ส Cloud & IT Cert Prep ให้อัปเกรดเป็น CoddyKit PRO คอร์ส Cloud & IT Cert Prep มีบทเรียนทั้งหมด 4 บทเรียน
คุณจะเรียนรู้อะไรในบทเรียน “Kinesis Streams, Firehose และการวิเคราะห์แบบเรียลไทม์”
รับข้อมูลสตรีมด้วย Kinesis Data Streams ส่งข้อมูลไปยัง S3 หรือ Redshift ด้วย Firehose และวิเคราะห์แบบเรียลไทม์ด้วย Managed Service for Apache Flink คุณปฏิบัติ Cloud & IT Cert Prep ด้วยโค้ดที่ใช้งานได้จริงที่คุณเรียกใช้โดยตรงในเบราว์เซอร์ และติวเตอร์ AI ตลอด 24/7 ตอบคำถามของคุณขณะที่คุณไปผ่านบทเรียน
คุณต้องมีประสบการณ์ก่อนที่จะเริ่มเรียน Cloud & IT Cert Prep หรือไม่
ไม่จำเป็นต้องมีประสบการณ์มาก่อน Cloud & IT Cert Prep บน CoddyKit ออกแบบมาสำหรับผู้เริ่มต้นไปจนถึงผู้เรียนขั้นสูง คุณสามารถเริ่มต้นที่นี่หรือเริ่มจากตัวแรกและเรียนด้วยความเร็วของคุณเอง นี่คือบทเรียนที่ 4 จากทั้งหมด 4 บทเรียน
บทเรียน “Kinesis Streams, Firehose และการวิเคราะห์แบบเรียลไทม์” ใช้เวลานานแค่ไหน
บทเรียน CoddyKit ส่วนใหญ่ใช้เวลาประมาณ 5–10 นาที แต่ละบทเรียนจึงสั้นและเป็นแบบโต้ตอบ คุณสามารถก้าวหน้าอย่างต่อเนื่องและกลับมาเรียนต่อจากตรงที่เพิ่งหยุดบนเว็บและแอปได้เลย
ฉันเขียนและรันโค้ดในบทเรียน Cloud & IT Cert Prep นี้ได้ไหม
ได้ บทเรียน Cloud & IT Cert Prep ทุกบทมีตัวแก้ไขโค้ดในตัว คุณจึงเขียนและรันโค้ดจริงได้เลยในเบราว์เซอร์ และได้รับข้อเสนอแนะจาก AI ในทันที — ไม่ต้องติดตั้งในเครื่องของคุณ
บทเรียนทั้งหมดในหลักสูตรนี้
- การสร้างดาต้าเลคบน S3
- AWS Glue: ETL และแค็ตตาล็อกข้อมูล
- Amazon Athena: SQL แบบไร้เซิร์ฟเวอร์บน S3
- Kinesis Streams, Firehose และการวิเคราะห์แบบเรียลไทม์