You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何将S3 Bucket数据导入Kinesis Stream?新手技术求助

How to Sync S3 Bucket Data to Amazon Kinesis Stream as a New User

Hey there! Let's walk through this clearly since you're new to Kinesis—we'll get your setup working smoothly.

First, let's clarify your question about Lambda's Kinesis events: that's not the right fit for your current use case. The Kinesis event trigger in Lambda is for when you want Lambda to consume data from a Kinesis Stream (i.e., read records out of Kinesis to process them). What you need here is the opposite: trigger Lambda when a new file lands in S3, then have Lambda send that file's data into your Kinesis Stream (mystream).

Here's the step-by-step implementation you should follow:

  1. Set up an S3 Event Trigger for your Lambda

    • Go to your Lambda function in the AWS Console, navigate to the "Configuration" tab > "Triggers" > "Add trigger"
    • Select "S3" as the trigger type, choose your target bucket
    • Under "Event types", pick All object create events (or specifically PutObject since you're uploading files)
    • Optional: Add a prefix/suffix filter if you only want to trigger on specific files (e.g., .csv or a specific folder)
    • Save the trigger—now your Lambda will run automatically whenever a new file is uploaded to that S3 bucket.
  2. Update your Lambda code to push data to Kinesis
    You already have code to read the S3 file and save to RDS—just add the Kinesis push logic. Here's a quick Python example using boto3:

    import boto3
    import json
    
    s3 = boto3.client('s3')
    kinesis = boto3.client('kinesis')
    
    def lambda_handler(event, context):
        # Get S3 file details from the trigger event
        bucket = event['Records'][0]['s3']['bucket']['name']
        key = event['Records'][0]['s3']['object']['key']
        
        # Read the file content from S3
        response = s3.get_object(Bucket=bucket, Key=key)
        file_content = response['Body'].read().decode('utf-8')
        
        # Optional: Parse the content if it's structured (e.g., JSON)
        # data = json.loads(file_content)
        
        # Push the data to Kinesis Stream
        kinesis_response = kinesis.put_record(
            StreamName='mystream',
            Data=file_content,  # Or the parsed data as JSON string
            PartitionKey='your-partition-key'  # Use a meaningful key (e.g., file name, record ID)
        )
        
        # Your existing RDS save logic goes here...
        
        return {
            'statusCode': 200,
            'body': json.dumps('Data pushed to Kinesis and saved to RDS successfully!')
        }
    
  3. Ensure your Lambda has the right IAM permissions
    Your Lambda's execution role needs:

    • Permissions to read from your S3 bucket (s3:GetObject for the bucket ARN)
    • Permissions to write to your Kinesis Stream (kinesis:PutRecord for the stream ARN)
    • The existing permissions you already have for RDS access

    You can add these to the role's policy in IAM Console—here's a sample policy snippet:

    {
        "Version": "2012-10-17",
        "Statement": [
            {
                "Effect": "Allow",
                "Action": "s3:GetObject",
                "Resource": "arn:aws:s3:::your-bucket-name/*"
            },
            {
                "Effect": "Allow",
                "Action": "kinesis:PutRecord",
                "Resource": "arn:aws:kinesis:your-region:your-account-id:stream/mystream"
            }
        ]
    }
    
  4. Test the flow

    • Upload a test file to your S3 bucket
    • Check the Lambda's "Monitor" tab > "Logs" to confirm it ran without errors
    • Go to your Kinesis Stream in the Console, use the "Data viewer" to check if the record was added

Quick note on your existing setup:

Since you're only uploading a record every 3 months, this setup is perfect—Lambda's pay-per-use model means you'll barely incur any costs, and it's fully automated so you don't have to manually trigger anything.

If you run into any snags (like permission errors, data formatting issues), feel free to share specific details and we can troubleshoot further!

内容的提问来源于stack exchange,提问作者ShaiNe Ram

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.15 07:43:29