Kinesis Streams وFirehose والتحليلات الفورية
استوعب البيانات المتدفقة باستخدام Kinesis Data Streams، وسلّمها إلى S3 أو Redshift باستخدام Firehose، وحلّلها في الوقت الفعلي باستخدام Managed Service for Apache Flink
Kinesis Streams وFirehose والتحليلات الفورية درس مجاني في Cloud & IT Cert Prep على CoddyKit. هذا هو الدرس 4 من أصل 4. يمكنك قراءة الدرس كاملاً أدناه مجاناً — ثم تمرن عليه مباشرة في المتصفح باستخدام محرر أكواد مدمج ومدرس ذكاء اصطناعي متاح 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 سجل دائم ومرتب، مقسم إلى 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-appKinesis 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) وفتح باقي دورة 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 مع أكواد عملية تشغلها مباشرة في المتصفح، ومدرس ذكاء اصطناعي متاح 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 يتضمن محرر أكواد مدمج، لذا تكتب وتشغل أكواداً حقيقية مباشرة في متصفحك وتحصل على تعليقات فورية من الذكاء الاصطناعي — بدون إعداد محلي.
جميع الدروس في هذه الدورة
- إنشاء بحيرة بيانات على S3
- AWS Glue: ETL وكتالوج البيانات
- Amazon Athena: SQL دون خوادم على S3
- Kinesis Streams وFirehose والتحليلات الفورية