Kinesis Streams, Firehose और रीयल-टाइम विश्लेषण
Kinesis Data Streams से स्ट्रीमिंग डेटा ग्रहण कीजिए, Firehose से उसे S3 या Redshift में भेजिए और Managed Service for Apache Flink से रीयल टाइम में विश्लेषण कीजिए।
Kinesis Streams, Firehose और रीयल-टाइम विश्लेषण, CoddyKit पर Cloud & IT Cert Prep का एक निःशुल्क पाठ है। यह 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 तक पूरी तरह प्रबंधित डिलीवरी), और रीयल-टाइम SQL तथा Flink प्रोसेसिंग के लिए Managed Service for Apache Flink (पहले Kinesis Data Analytics)। प्रत्येक सेवा स्ट्रीमिंग pipeline के अलग उद्देश्य को पूरा करती है।
Kinesis Data Streams आर्किटेक्चर
Kinesis Data Stream एक टिकाऊ, क्रमबद्ध लॉग है, जिसे shards में विभाजित किया गया है। प्रत्येक shard 1 MB/s लेखन थ्रूपुट और 2 MB/s पठन थ्रूपुट प्रदान करता है। डेटा रिकॉर्ड डिफ़ॉल्ट रूप से 24 घंटे तक रखे जाते हैं (इसे 7 दिन या 365 दिन तक बढ़ाया जा सकता है)। Producer किसी stream में रिकॉर्ड लिखते हैं; consumer — Lambda, KCL applications, Firehose या Flink — एक या अधिक shards से समानांतर रूप से पढ़ते हैं। लिखे जाने के बाद रिकॉर्ड बदले नहीं जा सकते।
# 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 की संख्या stream का कुल थ्रूपुट निर्धारित करती है। थ्रूपुट दोगुना करने के लिए आप एक shard को विभाजित कर सकते हैं या लागत घटाने के लिए दो shards को मिलाकर एक कर सकते हैं। Enhanced Fan-Out का उपयोग करके प्रत्येक पंजीकृत consumer को अन्य consumers से स्वतंत्र 2 MB/s का अपना पठन थ्रूपुट दें। इससे कई applications द्वारा एक ही stream का उपभोग करने पर read throttling समाप्त हो जाती है। Consumer की देरी का पता लगाने के लिए 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 endpoints सहित विभिन्न destinations तक पहुँचाती है। प्रबंधित करने के लिए कोई shards नहीं होते—Firehose अपने-आप scale करता है। आप buffer size (1–128 MB) और buffer interval (60–900 seconds) कॉन्फ़िगर करते हैं; दोनों में से कोई भी सीमा पहले पूरी होने पर 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"
}'Lambda के साथ Firehose डेटा रूपांतरण
Firehose डिलीवरी से पहले डेटा को रूपांतरित, समृद्ध या फ़िल्टर करने के लिए रिकॉर्ड के प्रत्येक batch पर Lambda function चला सकता है। सामान्य उपयोगों में JSON को Parquet में बदलना (Glue schema के माध्यम से), PII फ़ील्ड को छिपाना या कम-मूल्य वाली घटनाओं को हटाना शामिल हैं। रूपांतरण में विफल रहने वाले रिकॉर्ड वैकल्पिक रूप से पुनः प्रोसेसिंग के लिए अलग S3 error prefix में भेजे जाते हैं, इसलिए कोई डेटा नष्ट नहीं होता।
# 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 के लिए Managed Service
Amazon Managed Service for Apache Flink (पहले Kinesis Data Analytics) पूरी तरह प्रबंधित इन्फ्रास्ट्रक्चर पर Apache Flink applications चलाता है। इसका उपयोग stateful, रीयल-टाइम विश्लेषण के लिए करें: sliding window aggregations, anomaly detection, event sequences पर pattern matching और streaming data को reference tables के साथ JOIN करना। आप Flink code Java, Python या Scala में लिखते हैं और Flink checkpoints तथा exactly-once state को प्रबंधित करता है।
# 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 परीक्षा के लिए निर्णय के मानदंड जानें। जब आपको एक सेकंड से कम विलंबता, एक साथ पढ़ने वाले कई consumers या retention और replay पर पूर्ण नियंत्रण के साथ कस्टम प्रोसेसिंग logic चाहिए, तब Kinesis Data Streams का उपयोग करें। जब आपको कम code, वैकल्पिक in-flight transformation और अधिक विलंबता (60+ seconds) पर automatic scaling के साथ स्ट्रीमिंग डेटा को विश्वसनीय रूप से S3, Redshift या OpenSearch तक पहुँचाना हो, तब Firehose का उपयोग करें।
Kinesis बनाम SQS: परीक्षा में मिलने वाला पारंपरिक चुनाव
SAA-C03 में अक्सर Kinesis और SQS के बीच चुनाव करने वाला प्रश्न पूछा जाता है। मुख्य अंतर: Kinesis एक shard के भीतर संदेशों का क्रम बनाए रखता है, कई consumers को एक ही डेटा एक साथ पढ़ने देता है और replay के लिए रिकॉर्ड रखता है। SQS संदेशों के उपभोग के बाद उन्हें हटा देता है (कोई replay नहीं), FIFO mode सख्त क्रम की गारंटी देता है और यह microservices को अलग रखने के लिए बेहतर है। यदि परिदृश्य में रीयल-टाइम विश्लेषण या replay का उल्लेख हो, तो Kinesis चुनें।
Kinesis Producers: SDK और KPL
Kinesis Producer Library (KPL) applications से Kinesis Data Streams में लिखने वाला उच्च-थ्रूपुट client है। KPL कई छोटे records को अपने-आप एक API call (अधिकतम 1 MB) में एकत्रित करता है और back-off के साथ retries संभालता है। इससे प्रति-record PUT लागत बहुत घटती है और प्रत्येक shard का थ्रूपुट बढ़ता है। web clickstreams, IoT sensors या log pipelines जैसे उच्च-मात्रा वाले producers के लिए 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 पैटर्न
एक सामान्य आर्किटेक्चर में Kinesis streams या सीधे producers से डेटा वेयरहाउसिंग के लिए Amazon Redshift तक प्रबंधित pipeline के रूप में Firehose का उपयोग किया जाता है। Firehose पहले डेटा को मध्यवर्ती S3 staging bucket में लिखता है, फिर डेटा को Redshift में लोड करने के लिए COPY command जारी करता है। Redshift में स्ट्रीमिंग डेटा को bulk-load करने का यह सबसे प्रभावी तरीका है—Redshift में सीधे पंक्ति-दर-पंक्ति inserts करना प्रत्येक पंक्ति के अतिरिक्त overhead के कारण बहुत धीमा होगा।
# 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"
}
}'CloudWatch के साथ Kinesis की निगरानी
Kinesis के लिए प्रमुख CloudWatch metrics: producer थ्रूपुट मापने के लिए IncomingBytes और IncomingRecords, consumer lag मापने के लिए GetRecords.IteratorAgeMilliseconds (अधिक मान का अर्थ है कि consumers डेटा की गति के साथ नहीं चल पा रहे हैं), और producers के shard सीमाओं तक पहुँचने का पता लगाने के लिए WriteProvisionedThroughputExceeded। iterator age और provisioned throughput exceeded पर CloudWatch alarms सेट करें, ताकि Application Auto Scaling के माध्यम से shards की 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 कई consumers और replay क्षमता के साथ टिकाऊ, क्रमबद्ध और shards में विभाजित स्ट्रीमिंग प्रदान करता है, Kinesis Data Firehose S3, Redshift और OpenSearch तक पूरी तरह प्रबंधित, बिना code वाली डिलीवरी देता है, और Managed Service for Apache Flink streams पर stateful रीयल-टाइम विश्लेषण सक्षम करता है। आगे हम cloud migration strategy के 7 Rs का अध्ययन करेंगे।
एआई शिक्षक के साथ Cloud & IT Cert Prep सीखें — निःशुल्क
अपने ब्राउज़र में वास्तविक कोड लिखें और चलाएँ, चौबीसों घंटे एआई शिक्षक से तुरंत सहायता पाएँ, और वेब या ऐप पर वहीं से शुरू करें जहाँ आपने छोड़ा था।
- पाठ्यक्रम
- 150
- पाठ
- 600
अक्सर पूछे जाने वाले प्रश्न
क्या “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 से स्ट्रीमिंग डेटा ग्रहण कीजिए, Firehose से उसे S3 या Redshift में भेजिए और Managed Service for Apache Flink से रीयल टाइम में विश्लेषण कीजिए। आप ब्राउज़र में सीधे चलाए जाने वाले व्यावहारिक कोड के साथ Cloud & IT Cert Prep का अभ्यास करते हैं, और पाठ पूरा करते समय 24/7 एआई ट्यूटर आपके प्रश्नों के उत्तर देता है।
क्या Cloud & IT Cert Prep शुरू करने के लिए मुझे किसी अनुभव की आवश्यकता है?
पहले के अनुभव की आवश्यकता नहीं है। CoddyKit पर Cloud & IT Cert Prep शुरुआती से लेकर उन्नत शिक्षार्थियों तक सभी के लिए व्यवस्थित किया गया है, इसलिए आप यहीं से या शुरुआत से सीखना शुरू कर सकते हैं और अपनी गति से आगे बढ़ सकते हैं। यह 4 में से 4वाँ पाठ है।
“Kinesis Streams, Firehose और रीयल-टाइम विश्लेषण” पाठ पूरा करने में कितना समय लगता है?
CoddyKit का अधिकांश पाठ लगभग 5–10 मिनट में पूरा हो जाता है। हर पाठ छोटा और संवादात्मक है, इसलिए आप लगातार प्रगति करते हैं और वेब या ऐप पर वहीं से सीखना जारी रख सकते हैं जहाँ आपने छोड़ा था।
क्या मैं इस Cloud & IT Cert Prep पाठ में कोड लिख और चला सकता हूँ?
हाँ। हर Cloud & IT Cert Prep पाठ में एक अंतर्निर्मित कोड संपादक शामिल है, जिससे आप सीधे अपने ब्राउज़र में वास्तविक कोड लिख और चला सकते हैं और तुरंत एआई प्रतिक्रिया पा सकते हैं—स्थानीय सेटअप की आवश्यकता नहीं है।
इस पाठ्यक्रम के सभी पाठ
- S3 पर डेटा लेक बनाना
- AWS Glue: ETL और डेटा कैटलॉग
- Amazon Athena: S3 पर सर्वररहित SQL
- Kinesis Streams, Firehose और रीयल-टाइम विश्लेषण