Kinesis Data Streams لمعالجة الأحداث الفورية
أنشئ واستهلك تدفقات أحداث عالية الإنتاجية باستخدام Kinesis Data Streams، وأدر الأجزاء لتحقيق معدل النقل، واستخدم Lambda كمستهلك
Kinesis Data Streams لمعالجة الأحداث الفورية درس مجاني في AWS Solutions Architect على CoddyKit. هذا هو الدرس 3 من أصل 4. يمكنك قراءة الدرس كاملاً أدناه مجاناً — ثم تمرن عليه مباشرة في المتصفح باستخدام محرر أكواد مدمج ومدرس ذكاء اصطناعي متاح 24/7. هذا الدرس جزء من مسار التعلم في AWS Solutions Architect، وتقدمك يتزامن عبر الويب وتطبيق CoddyKit. تتضمن دورة AWS Solutions Architect 4 دروس في المجموع.
المفاهيم الأساسية في Kinesis Data Streams
إن Kinesis Data Streams (KDS) خدمة متينة ومرتبة لبث البيانات في الوقت الفعلي. تُنظَّم البيانات في stream يتكون من shards واحد أو أكثر. وكل shard عبارة عن تسلسل مرتب من سجلات البيانات. يضع المنتجون السجلات في shards باستخدام مفتاح التقسيم الذي يحدد shard الذي يستقبل السجل. ويقرأ المستهلكون السجلات من shards ويعالجونها بالترتيب الذي وصلت به ضمن كل shard.
سعة Shard وحدود الإنتاجية
يدعم كل shard إنتاجية كتابة تبلغ 1 ميغابايت/ثانية أو 1,000 سجل/ثانية وإنتاجية قراءة تبلغ 2 ميغابايت/ثانية (تتشارك فيها جميع المستهلكين القياسيين على ذلك shard). وتتوسع السعة الإجمالية للتدفق خطيًا مع عدد shards. استخدم الصيغة: shards_needed = max(write_MB_per_s / 1, read_MB_per_s / 2). إذا بلغ المنتجون حد الكتابة، فستظهر أخطاء ProvisionedThroughputExceededException — ويمكن حل ذلك بتقسيم shards أو توزيع مفاتيح التقسيم بالتساوي.
# Calculate shards needed for a stream:
# - Ingest rate: 5 MB/s writes
# - Read rate: 3 consumers x 2 MB/s = 6 MB/s reads
# shards = max(5/1, 6/2) = max(5, 3) = 5 shards needed
aws kinesis create-stream \
--stream-name iot-telemetry \
--shard-count 5مفاتيح التقسيم وتوزيع البيانات
مفتاح التقسيم عبارة عن سلسلة نصية تُجري Kinesis تجزئة لها باستخدام (MD5) لتحديد shard الذي يستقبل السجل. ويوزع مفتاح التقسيم المختار جيدًا السجلات بالتساوي عبر shards (منع امتلاء shard). بالنسبة إلى إنترنت الأشياء، استخدم معرّف الجهاز. وبالنسبة إلى تدفقات النقرات، استخدم معرّف الجلسة أو معرّف المستخدم. تجنب مفاتيح التقسيم منخفضة التنوع (مثل اسم الدولة الذي لا يتضمن سوى 5 قيم)، لأنها تؤدي إلى امتلاء shards، حيث يستقبل أحدها قدرًا غير متناسب من حركة الكتابة بينما تظل الأخرى خاملة.
import boto3, json
client = boto3.client('kinesis', region_name='us-east-1')
# Good: use device_id as partition key for even distribution
event = {'deviceId': 'sensor-42', 'temp': 23.5, 'ts': '2024-01-15T10:00:00Z'}
client.put_record(
StreamName='iot-telemetry',
Data=json.dumps(event),
PartitionKey='sensor-42' # high-cardinality -> even distribution
)المستهلكون القياسيون مقابل Enhanced Fan-Out
يتشارك المستهلكون القياسيون إنتاجية القراءة البالغة 2 ميغابايت/ثانية لكل shard، باستخدام GetRecords مع الاستقصاء. إذا كان لديك 3 مستهلكين على shard، ويحتاج كل منهم إلى 2 ميغابايت/ثانية، فسيُفرض عليهم تقييد المعدل أثناء مشاركتهم في إجمالي قدره 2 ميغابايت/ثانية. يوفر Enhanced Fan-Out (EFO) لكل مستهلك مسجل قناة قراءة مخصصة تبلغ 2 ميغابايت/ثانية، عبر اتصال دفع دائم باستخدام HTTP/2 (SubscribeToShard). يضيف EFO تكلفة لكل مستهلك ولكل ساعة shard، لكنه يلغي التنافس على القراءة بالكامل.
# Register an Enhanced Fan-Out consumer
aws kinesis register-stream-consumer \
--stream-arn arn:aws:kinesis:us-east-1:123456789012:stream/iot-telemetry \
--consumer-name real-time-analytics
# List registered consumers
aws kinesis list-stream-consumers \
--stream-arn arn:aws:kinesis:us-east-1:123456789012:stream/iot-telemetryLambda كمستهلك لـ Kinesis
يتكامل Lambda أصليًا مع Kinesis Data Streams عبر Event Source Mapping. يستقصي Lambda التدفق، ويقرأ دفعات من السجلات، ثم يستدعي دالتك. اضبط BatchSize (من 1 إلى 10,000 سجل)، وStartingPosition (TRIM_HORIZON للأقدم وLATEST للأحدث)، وBisectBatchOnFunctionError لتقسيم الدفعات الفاشلة. ويتيح Parallelisation Factor (من 1 إلى 10) لـ Lambda تشغيل استدعاءات متزامنة متعددة لكل shard لمواكبة التدفقات السريعة.
# Create a Lambda event source mapping for Kinesis
aws lambda create-event-source-mapping \
--function-name ProcessIoTEvents \
--event-source-arn arn:aws:kinesis:us-east-1:123456789012:stream/iot-telemetry \
--batch-size 100 \
--starting-position LATEST \
--parallelization-factor 5 \
--bisect-batch-on-function-error true \
--destination-config '{
"OnFailure": {
"Destination": "arn:aws:sqs:us-east-1:123456789012:kinesis-dlq"
}
}'مكتبة عميل Kinesis (KCL)
إن Kinesis Client Library (KCL) إطار عمل لبناء مستهلكي Kinesis متينين باستخدام Java (أو لغات متعددة عبر MultiLangDaemon). تتولى KCL تعداد shards، وإدارة عقود الإيجار (توزيع shards بين مثيلات العاملين)، وحفظ نقاط التحقق الخاصة بالتقدم في DynamoDB، والتعامل السلس مع تقسيم shards ودمجها. يعالج كل عامل من عمال KCL shard واحدًا أو أكثر، وتعيد KCL موازنة shards تلقائيًا عند زيادة عدد العاملين أو تعطلهم. ويُفضَّل استخدام KCL على الاستقصاء المباشر عبر SDK في تطبيقات المستهلكين الإنتاجية.
# KCL stores checkpoints in a DynamoDB table automatically
# Each shard has one row tracking the last successfully processed sequence number
# KCL lease table structure:
# leaseKey (shardId) | checkpoint (sequenceNumber) | leaseOwner (workerId)
#
# To start a KCL application (pseudocode):
# KinesisClientLibConfiguration config = new KinesisClientLibConfiguration(
# 'iot-app', 'iot-telemetry', credentialsProvider, 'worker-1');
# Worker worker = new Worker.Builder().config(config).recordProcessorFactory(factory).build();
# worker.run();الاحتفاظ بالبيانات وإعادة تشغيلها
تخزن Kinesis Data Streams السجلات لمدة 24 ساعة افتراضيًا (ويمكن تمديدها إلى 7 أيام أو حتى 365 يومًا باستخدام الاحتفاظ طويل الأمد مقابل تكلفة إضافية). بخلاف SQS، لا تُحذف السجلات المقروءة بعد استهلاكها، بل تظل متاحة حتى انتهاء مدة الاحتفاظ. ويتيح ذلك لمستهلكين متعددين قراءة السجلات نفسها بشكل مستقل، كما يتيح إعادة التشغيل عبر إعادة تعيين نقطة تحقق المستهلك إلى رقم تسلسل أسبق، وهو أمر بالغ القيمة لإصلاح الأخطاء أو تعبئة الخدمات الجديدة بالبيانات السابقة.
# Extend stream retention to 7 days
aws kinesis increase-stream-retention-period \
--stream-name iot-telemetry \
--retention-period-hours 168
# Get records from the oldest available record (replay)
SHARD_ITERATOR=$(aws kinesis get-shard-iterator \
--stream-name iot-telemetry \
--shard-id shardId-000000000000 \
--shard-iterator-type TRIM_HORIZON \
--query 'ShardIterator' --output text)
aws kinesis get-records --shard-iterator $SHARD_ITERATOR --limit 100وضع السعة عند الطلب
تدعم Kinesis Data Streams وضعي سعة. الوضع المُوفَّر: تدير عدد shards يدويًا وتدفع التكلفة لكل shard ولكل ساعة. الوضع عند الطلب: توسّع Kinesis سعة shards تلقائيًا استنادًا إلى الإنتاجية الواردة (حتى 200 ميغابايت/ثانية للكتابة و400 ميغابايت/ثانية للقراءة افتراضيًا)، وتدفع التكلفة لكل غيغابايت من البيانات المكتوبة والمستردة. يناسب الوضع عند الطلب أنماط حركة المرور المتغيرة أو غير المتوقعة، عندما لا ترغب في إدارة توسعة shards.
# Switch an existing stream to On-Demand mode
aws kinesis update-stream-mode \
--stream-arn arn:aws:kinesis:us-east-1:123456789012:stream/iot-telemetry \
--stream-mode-details StreamMode=ON_DEMANDضمانات الترتيب داخل Shards
تضمن Kinesis الترتيب داخل shard — إذ تنتقل السجلات ذات مفتاح التقسيم نفسه دائمًا إلى shard نفسه، وتُقرأ بالترتيب الذي كُتبت به. ومع ذلك، لا يوجد ضمان للترتيب عبر shards. إذا كان تطبيقك يتطلب ترتيبًا شاملًا لجميع السجلات، فاستخدم shard واحدًا (مما يحد الإنتاجية إلى 1 ميغابايت/ثانية)، أو أعد التصميم بحيث لا تكون هناك حاجة إلى الترتيب إلا ضمن مجموعة مفتاح تقسيم (مثل الترتيب لكل جهاز). وهذا تمييز شائع في الاختبار مقارنةً بـ SQS FIFO، الذي يوفر إزالة تكرار وترتيبًا صارمين.
الأمان: التشفير وVPC
تشفّر Kinesis Data Streams السجلات من جهة الخادم أثناء السكون باستخدام AWS KMS (مفتاح CMK أو مفتاح تديره AWS) عند تمكين التشفير من جهة الخادم. وتُشفَّر جميع البيانات أثناء النقل باستخدام TLS. بالنسبة إلى التطبيقات التي تعمل في VPC ولا ينبغي لها إرسال بيانات التدفق عبر الإنترنت العام، استخدم VPC Interface Endpoint (PrivateLink) لـ Kinesis، بحيث تبقى حركة المرور بالكامل ضمن العمود الفقري لشبكة AWS — وهو أمر مهم لأحمال العمل الخاضعة للأنظمة الرقابية.
# Enable server-side encryption on a Kinesis stream
aws kinesis start-stream-encryption \
--stream-name iot-telemetry \
--encryption-type KMS \
--key-id arn:aws:kms:us-east-1:123456789012:key/mrk-abc123Kinesis Data Streams مقابل SQS وKafka
في اختبار SAA-C03، قارن بين Kinesis Data Streams والبدائل. Kinesis مقابل SQS: تحافظ Kinesis على الترتيب داخل shard وتدعم قراءة المستهلكين المتعددين للبيانات نفسها، بينما تحذف SQS الرسائل بعد استهلاكها. Kinesis مقابل MSK (Kafka): إن MSK هو Apache Kafka مُدار — استخدمه عندما تحتاج إلى توافق بروتوكول Kafka، أو إعدادات متقدمة للموضوعات، أو عند الترحيل من Kafka المحلي. استخدم Kinesis للبث الأصلي على AWS مع تكامل أقوى مع Lambda وFirehose وFlink. اختر Kinesis ما لم يذكر السؤال Kafka تحديدًا أو متطلبات التوافق معه.
تحقق سريع
اختبر مدى فهمك لمفاهيم AWS Solutions Architect (SAA-C03) الواردة في هذا الدرس.
مراجعة الدرس
تعلمت في هذا الدرس أن: Kinesis Data Streams توفر تدفقات مقسمة إلى shards ومرتبة ومتينة، مع احتفاظ افتراضي لمدة 24 ساعة وإمكانية إعادة التشغيل، وأن Enhanced Fan-Out يمنح كل مستهلك سرعة مخصصة تبلغ 2 ميغابايت/ثانية لكل shard، مما يلغي التنافس على القراءة، وأن الوضع عند الطلب يوسّع shards تلقائيًا للتعامل مع حركة المرور غير المتوقعة. سنتناول بعد ذلك نمطي التنسيق بالتتابع مقابل التنظيم في البنية القائمة على الأحداث.
الأسئلة الشائعة
هل درس «Kinesis Data Streams لمعالجة الأحداث الفورية» مجاني؟
نعم — نص درس «Kinesis Data Streams لمعالجة الأحداث الفورية» كامل متاح مجاناً هنا على الويب. لتمرينه بشكل تفاعلي (محرر أكواد مدمج ومدرس ذكاء اصطناعي متاح 24/7) وفتح باقي دورة AWS Solutions Architect، انتقل إلى CoddyKit PRO. تتضمن دورة AWS Solutions Architect 4 دروس في المجموع.
ماذا ستتعلم في «Kinesis Data Streams لمعالجة الأحداث الفورية»؟
أنشئ واستهلك تدفقات أحداث عالية الإنتاجية باستخدام Kinesis Data Streams، وأدر الأجزاء لتحقيق معدل النقل، واستخدم Lambda كمستهلك تتمرن على AWS Solutions Architect مع أكواد عملية تشغلها مباشرة في المتصفح، ومدرس ذكاء اصطناعي متاح 24/7 يجيب على أسئلتك أثناء عملك.
هل أحتاج إلى خبرة سابقة لأبدأ AWS Solutions Architect؟
لا تُشترط خبرة سابقة. AWS Solutions Architect على CoddyKit منظم للمبتدئين حتى المتقدمين، لذا يمكنك البدء من هنا أو من البداية والتقدم بسرعتك الخاصة. هذا هو الدرس 3 من أصل 4.
كم من الوقت يستغرق درس «Kinesis Data Streams لمعالجة الأحداث الفورية»؟
معظم دروس CoddyKit تستغرق حوالي 5–10 دقائق. كل منها موجز وتفاعلي، لذا تحرز تقدماً مستمراً وتستأنف من حيث توقفت عبر الويب والتطبيق.
هل يمكنني كتابة وتشغيل أكواد في درس AWS Solutions Architect هذا؟
نعم. كل درس في AWS Solutions Architect يتضمن محرر أكواد مدمج، لذا تكتب وتشغل أكواداً حقيقية مباشرة في متصفحك وتحصل على تعليقات فورية من الذكاء الاصطناعي — بدون إعداد محلي.
جميع الدروس في هذه الدورة
- EventBridge: ناقل الأحداث والقواعد
- Step Functions: تنسيق سير العمل دون خوادم
- Kinesis Data Streams لمعالجة الأحداث الفورية
- أنماط التنسيق مقابل التوجيه المركزي