0Pricing
AWS Solutions Architect · درس

Kinesis Streams وFirehose والتحليلات الفورية

استوعب البيانات المتدفقة باستخدام Kinesis Data Streams، وسلّمها إلى S3 أو Redshift باستخدام Firehose، وحلّلها في الوقت الفعلي باستخدام Managed Service for Apache Flink

Kinesis Streams وFirehose والتحليلات الفورية درس مجاني في AWS Solutions Architect على CoddyKit. هذا هو الدرس 4 من أصل 4. يمكنك قراءة الدرس كاملاً أدناه مجاناً — ثم تمرن عليه مباشرة في المتصفح باستخدام محرر أكواد مدمج ومدرس ذكاء اصطناعي متاح 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 سجل دائم ومرتب، مقسم إلى shards. يوفر كل shard معدل نقل يبلغ 1 MB/s للكتابة و2 MB/s للقراءة. وتُحتفظ بسجلات البيانات لمدة 24 ساعة افتراضيًا، ويمكن تمديدها إلى 7 أيام أو 365 يومًا. يكتب المنتجون السجلات إلى تدفق، بينما يقرأ المستهلكون — مثل Lambda وتطبيقات KCL وFirehose وFlink — من shard واحد أو أكثر بالتوازي. وتصبح السجلات غير قابلة للتغيير بعد كتابتها.

# 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'

Shards ومعدل النقل والتوسعة

يحدد عدد shards إجمالي معدل نقل التدفق. ويمكنكم تقسيم shard لمضاعفة معدل النقل، أو دمج shardين لخفض التكلفة. استخدموا 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. ولا توجد shards لإدارتها، إذ تتوسع 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}

Managed Service for Apache Flink

تشغّل Amazon Managed Service for Apache Flink (المعروفة سابقًا باسم Kinesis Data Analytics) تطبيقات Apache Flink على بنية تحتية مُدارة بالكامل. استخدموها للتحليلات ذات الحالة في الوقت الفعلي، مثل التجميع باستخدام النوافذ المنزلقة، واكتشاف الحالات الشاذة، ومطابقة الأنماط ضمن تسلسلات الأحداث، وربط بيانات البث بجداول مرجعية. تكتبون كود 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 على ترتيب الرسائل داخل shard، وتدعم قراءة عدة مستهلكين للبيانات نفسها في الوقت نفسه، وتحتفظ بالسجلات لإعادة تشغيلها. أما SQS فتزيل الرسائل بعد استهلاكها (من دون إمكانية إعادة التشغيل)، ويضمن وضع FIFO فيها ترتيبًا صارمًا، كما أنها أنسب لفصل الخدمات المصغرة. فإذا ذكر السيناريو تحليلات في الوقت الفعلي أو إعادة التشغيل، فاختاروا Kinesis.

منتجو Kinesis: SDK وKPL

إن Kinesis Producer Library (KPL) عميل عالي معدل النقل للكتابة إلى Kinesis Data Streams من التطبيقات. تعمل KPL تلقائيًا على تجميع عدة سجلات صغيرة في استدعاء API واحد (حتى 1 MB)، كما تتولى إعادة المحاولة مع التراجع التدريجي. ويؤدي ذلك إلى تقليل تكلفة PUT لكل سجل بشكل كبير وزيادة معدل النقل لكل shard. استخدموا KPL للمنتجين ذوي الأحجام الكبيرة، مثل تدفقات نقرات الويب، وأجهزة استشعار إنترنت الأشياء، ومسارات السجلات.

# 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 لاكتشاف وصول المنتجين إلى حدود shard. اضبطوا تنبيهات CloudWatch بشأن عمر المؤشر وتجاوز معدل النقل المخصص لتشغيل التوسعة التلقائية لـ shards عبر 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 بثًا دائمًا ومرتبًا ومقسمًا إلى shards، مع دعم عدة مستهلكين وإمكانية إعادة التشغيل، وتوفر Kinesis Data Firehose تسليمًا مُدارًا بالكامل ومن دون تعليمات برمجية إلى S3 وRedshift وOpenSearch، وتتيح Managed Service for Apache Flink إجراء تحليلات ذات حالة في الوقت الفعلي على تدفقات البيانات. بعد ذلك، سنستكشف الاستراتيجية السحابية لترحيل البيانات وفق قواعد Rs السبع.

الأسئلة الشائعة

هل درس «Kinesis Streams وFirehose والتحليلات الفورية» مجاني؟

نعم — نص درس «Kinesis Streams وFirehose والتحليلات الفورية» كامل متاح مجاناً هنا على الويب. لتمرينه بشكل تفاعلي (محرر أكواد مدمج ومدرس ذكاء اصطناعي متاح 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 مع أكواد عملية تشغلها مباشرة في المتصفح، ومدرس ذكاء اصطناعي متاح 24/7 يجيب على أسئلتك أثناء عملك.

هل أحتاج إلى خبرة سابقة لأبدأ AWS Solutions Architect؟

لا تُشترط خبرة سابقة. AWS Solutions Architect على CoddyKit منظم للمبتدئين حتى المتقدمين، لذا يمكنك البدء من هنا أو من البداية والتقدم بسرعتك الخاصة. هذا هو الدرس 4 من أصل 4.

كم من الوقت يستغرق درس «Kinesis Streams وFirehose والتحليلات الفورية»؟

معظم دروس CoddyKit تستغرق حوالي 5–10 دقائق. كل منها موجز وتفاعلي، لذا تحرز تقدماً مستمراً وتستأنف من حيث توقفت عبر الويب والتطبيق.

هل يمكنني كتابة وتشغيل أكواد في درس AWS Solutions Architect هذا؟

نعم. كل درس في AWS Solutions Architect يتضمن محرر أكواد مدمج، لذا تكتب وتشغل أكواداً حقيقية مباشرة في متصفحك وتحصل على تعليقات فورية من الذكاء الاصطناعي — بدون إعداد محلي.

جميع الدروس في هذه الدورة

  1. إنشاء بحيرة بيانات على S3
  2. AWS Glue: ETL وكتالوج البيانات
  3. Amazon Athena: SQL دون خوادم على S3
  4. Kinesis Streams وFirehose والتحليلات الفورية
← العودة إلى AWS Solutions Architect