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

Flask+Boto3操作DynamoDB数据获取不一致问题求助

问题描述

已配置AWS IoT规则将MQTT消息存储至DynamoDB v2表,存储流程正常——Arduino设备发送的MQTT数据能成功写入DynamoDB。但EC2上部署的Flask应用获取数据时出现异常:

  • AWS控制台扫描可见58条数据(默认显示50条,需点击“获取下一页”)
  • Flask应用返回数据量不稳定,有时50条、有时54条,偶尔才返回全部58条
  • 重启apache2服务器后,数据计数恢复正确

尝试两种实现方式均存在相同问题:

  1. 手动分页执行scan操作获取全量数据
  2. 使用boto3的paginator执行query操作
    两种方式均为重启apache2后首次运行能获取正确数据量,但新增数据后,必须重启服务器才能显示新数据。

代码示例1(手动分页Scan)

from flask import Flask
import boto3
from boto3.dynamodb.conditions import Key, Attr
from datetime import datetime

AWS_ACCESS_KEY_ID = 'key'
AWS_SECRET_ACCESS_KEY = 'access_key'
REGION_NAME = 'region'

dynamodb = boto3.resource(
    'dynamodb',
    aws_access_key_id     = AWS_ACCESS_KEY_ID,
    aws_secret_access_key = AWS_SECRET_ACCESS_KEY,
    region_name           = REGION_NAME
)

table = dynamodb.Table('my_db')

lastEvaluatedKey = None
items = []

#I'm doing this below to get all pages
while True:
    if lastEvaluatedKey == None:
        response = table.scan()
    else:
        response = table.scan(
        ExclusiveStartKey=lastEvaluatedKey
    )

    items.extend(response['Items'])

    if 'LastEvaluatedKey' in response:
        lastEvaluatedKey = response['LastEvaluatedKey']
    else:
        break

#localDateTime is my sortKey
sorted_items = sorted(items, key=lambda x: datetime.strptime(x['localDateTime'], '%Y-%m-%dT%H:%M:%S'), reverse=True)

app = Flask(__name__)
@app.route('/')

def my_app():
    return sorted_items
if __name__ == '__main__':
  app.run()

代码示例2(使用Paginator Query)

from flask import Flask
import boto3
from boto3.dynamodb.conditions import Key, Attr
from datetime import datetime

AWS_ACCESS_KEY_ID = 'key'
AWS_SECRET_ACCESS_KEY = 'secret'
REGION_NAME = 'region'

dynamodb = boto3.client(
    'dynamodb',
    aws_access_key_id     = AWS_ACCESS_KEY_ID,
    aws_secret_access_key = AWS_SECRET_ACCESS_KEY,
    region_name           = REGION_NAME
)

paginator = dynamodb.get_paginator('query')
response = paginator.paginate(TableName='my-table',
        KeyConditionExpression="id = :id",
        ExpressionAttributeValues= {
            ":id": {
                "S": "my_id"
                }
            }
        )

items = []

for page in response:
    items.append(page['Items'])


app = Flask(__name__)
@app.route('/')

def my_app():
    return items[0]
if __name__ == '__main__':
  app.run()
排查思路与解决方案

核心原因分析

问题根源在于数据查询逻辑被放在了Flask应用的全局作用域,而非请求处理函数内部。当Flask应用被apache2托管时,服务器启动时会加载一次全局代码,查询数据并缓存items或sorted_items变量。后续所有请求都复用这个缓存的旧数据,只有重启服务器才会重新执行查询获取最新数据。

同时,数据量不稳定的情况是因为DynamoDB的scan/query操作存在最终一致性特性,首次查询时可能未同步所有最新写入的数据,但重启后重新查询会拿到更全的数据,但后续新增数据依然不会自动更新。

解决方案

1. 将查询逻辑移至请求处理函数内部

每次请求时都重新执行DynamoDB查询,确保获取最新数据:

修改后的代码示例1(手动分页Scan):

from flask import Flask
import boto3
from boto3.dynamodb.conditions import Key, Attr
from datetime import datetime

AWS_ACCESS_KEY_ID = 'key'
AWS_SECRET_ACCESS_KEY = 'access_key'
REGION_NAME = 'region'

dynamodb = boto3.resource(
    'dynamodb',
    aws_access_key_id     = AWS_ACCESS_KEY_ID,
    aws_secret_access_key = AWS_SECRET_ACCESS_KEY,
    region_name           = REGION_NAME
)

table = dynamodb.Table('my_db')

app = Flask(__name__)
@app.route('/')
def my_app():
    lastEvaluatedKey = None
    items = []
    # 每次请求时重新执行分页查询
    while True:
        if lastEvaluatedKey is None:
            response = table.scan()
        else:
            response = table.scan(ExclusiveStartKey=lastEvaluatedKey)
        
        items.extend(response['Items'])
        
        if 'LastEvaluatedKey' in response:
            lastEvaluatedKey = response['LastEvaluatedKey']
        else:
            break
    
    sorted_items = sorted(items, key=lambda x: datetime.strptime(x['localDateTime'], '%Y-%m-%dT%H:%M:%S'), reverse=True)
    return sorted_items

if __name__ == '__main__':
  app.run()

修改后的代码示例2(使用Paginator Query):

from flask import Flask
import boto3
from boto3.dynamodb.conditions import Key, Attr
from datetime import datetime

AWS_ACCESS_KEY_ID = 'key'
AWS_SECRET_ACCESS_KEY = 'secret'
REGION_NAME = 'region'

dynamodb = boto3.client(
    'dynamodb',
    aws_access_key_id     = AWS_ACCESS_KEY_ID,
    aws_secret_access_key = AWS_SECRET_ACCESS_KEY,
    region_name           = REGION_NAME
)

app = Flask(__name__)
@app.route('/')
def my_app():
    paginator = dynamodb.get_paginator('query')
    response = paginator.paginate(TableName='my-table',
            KeyConditionExpression="id = :id",
            ExpressionAttributeValues= {
                ":id": {
                    "S": "my_id"
                    }
                }
            )
    
    items = []
    for page in response:
        items.extend(page['Items'])  # 改为extend避免嵌套列表
    return items

if __name__ == '__main__':
  app.run()

2. 处理DynamoDB最终一致性问题(可选)

如果需要强一致性读取,可以在query时添加ConsistentRead=True参数(scan操作不支持强一致性):

response = paginator.paginate(
    TableName='my-table',
    KeyConditionExpression="id = :id",
    ExpressionAttributeValues={":id": {"S": "my_id"}},
    ConsistentRead=True
)

注意:强一致性读取会增加延迟和消耗更多容量单位,需根据业务需求权衡。

3. 优化查询性能(可选)

如果数据量较大,每次请求全量查询会影响性能,可以考虑:

  • 使用DynamoDB的时间范围查询,仅返回最近一段时间的数据
  • 实现本地缓存策略,设置合理的缓存过期时间,避免频繁查询数据库

内容的提问来源于stack exchange,提问作者Rafael Cruz

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 15:35:55