0Pricing
Cloud & IT Cert Prep · บทเรียน

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-app

Kinesis 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 ในทันที — ไม่ต้องติดตั้งในเครื่องของคุณ

บทเรียนทั้งหมดในหลักสูตรนี้

  1. การสร้างดาต้าเลคบน S3
  2. AWS Glue: ETL และแค็ตตาล็อกข้อมูล
  3. Amazon Athena: SQL แบบไร้เซิร์ฟเวอร์บน S3
  4. Kinesis Streams, Firehose และการวิเคราะห์แบบเรียลไทม์
← กลับไปที่ Cloud & IT Cert Prep