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

如何在Android AWS Amplify中监听DynamoDB后端数据变更?

问题描述

我基于Android Amplify开发了AWS Amplify移动端应用,需求是实时监听DynamoDB中的数据变更——包括通过其他应用脚本或控制台直接修改的数据。目前通过移动端客户端添加数据时,变更数据能返回至我的监听方法,但当有人在控制台直接修改或编辑数据时,我无法接收到变更通知。

相关代码如下:

package com.example.appsyncdynamodb

import android.os.Bundle
import android.os.Handler
import android.os.Looper
import android.util.Log
import android.view.View
import android.widget.Button
import android.widget.Toast
import androidx.activity.enableEdgeToEdge
import androidx.appcompat.app.AppCompatActivity
import androidx.core.view.ViewCompat
import androidx.core.view.WindowInsetsCompat
import com.amplifyframework.api.graphql.model.ModelQuery
import com.amplifyframework.api.graphql.model.ModelSubscription
import com.amplifyframework.core.Action
import com.amplifyframework.core.Amplify
import com.amplifyframework.core.Consumer
import com.amplifyframework.core.async.Cancelable
import com.amplifyframework.core.model.query.ObserveQueryOptions
import com.amplifyframework.core.model.query.predicate.QueryPredicate
import com.amplifyframework.datastore.DataStoreException
import com.amplifyframework.datastore.DataStoreItemChange
import com.amplifyframework.datastore.DataStoreQuerySnapshot
import com.amplifyframework.datastore.generated.model.Todo

class MainActivity : AppCompatActivity() {
    override fun onCreate(savedInstanceState: Bundle?) {
        super.onCreate(savedInstanceState)
        enableEdgeToEdge()
        setContentView(R.layout.activity_main)
        ViewCompat.setOnApplyWindowInsetsListener(findViewById(R.id.main)) { v, insets ->
            val systemBars = insets.getInsets(WindowInsetsCompat.Type.systemBars())
            v.setPadding(systemBars.left, systemBars.top, systemBars.right, systemBars.bottom)
            insets
        }
        Amplify.DataStore.start({}, 
        {})
        val btn:Button = findViewById(R.id.btn)
        var i : Int=0
        btn.setOnClickListener {
            i += 1;
            Log.d("printed inside","inside")

            //your implementation goes here
            val item = Todo.builder()
                .name("Todo").description("Hello Vibhor")
                .build()

            Amplify.DataStore.save(item,
                { Log.i("Tutorial", "Saved item: ${item.toString()}") },
                { Log.e("Tutorial", "Could not save item to DataStore", it) }
            )
        }
        Log.d("printed outside","outside")


        val subscription = Amplify.API.subscribe(
            ModelSubscription.onUpdate(Todo::class.java),
            { Log.d("ApiQuickStart", "Subscription established") },
            { Log.d("ApiQuickStart", "Todo update subscription received:$it") },
            { Log.d("ApiQuickStart", "Subscription failed", it) },
            { Log.d("ApiQuickStart", "Subscription completed") }
        )

        Amplify.DataStore.observe(Todo::class.java,
            { Log.i("MyAmplifyApp", "Observation began") },
            {
                val post = it.item()
                Handler(Looper.getMainLooper()).post {
                    Toast.makeText(this, it.item().toString(), Toast.LENGTH_LONG).show()
                }
                Log.i("MyAmplifyAppobs", "Post: $post")
            },
            { Log.e("MyAmplifyApp", "Observation failed", it) },
            { Log.i("MyAmplifyApp", "Observation complete") }
        )

        val subscription2 = Amplify.API.subscribe(
            ModelSubscription.onCreate(Todo::class.java),
            { Log.d("ApiQuickStart", "Subscription established") },
            { Log.d("ApiQuickStart", "Todo create subscription received:$it") },
            { Log.d("ApiQuickStart", "Subscription failed", it) },
            { Log.d("ApiQuickStart", "Subscription completed") }
        )

        val tag = "ObserveQuery"
        val onQuerySnapshot: Consumer<DataStoreItemChange<Todo>> =
            Consumer<DataStoreItemChange<Todo>> { value: DataStoreItemChange<Todo> ->
                Log.d(tag, "success on snapshot")
                Log.d(tag, "records: $value")
                //Log.d(tag, "number of records: " + value.items.size)
                //Log.d(tag, "sync status: " + value.isSynced)
            }

        val observationStarted =
            Consumer { _: Cancelable ->
                Log.d(tag, "success on cancelable")
            }
        val onObservationError =
            Consumer { value: DataStoreException ->
                Log.d(tag, "error on snapshot $value")
            }
        val onObservationComplete = Action {
            Log.d(tag, "complete")
        }
       val predicate: QueryPredicate =
            Todo.NAME.beginsWith("Build")
//        val querySortBy = QuerySortBy("post", "rating", QuerySortOrder.ASCENDING)

//        val options = ObserveQueryOptions(predicate,null)
//        Amplify.DataStore.observeQuery(
//            Todo::class.java,
//            options,
//            observationStarted,
//            onQuerySnapshot,
//            onObservationError,
//            onObservationComplete
//        )


        Amplify.API.query(
            ModelQuery.list(Todo::class.java),
            { response ->
                response.data.forEach { todo ->
                    Log.i("Query++", todo.name)
                }
            },
            { Log.e("MyAmplifyApp", "Query failure", it) }
        )

        Amplify.DataStore.observe(
            {Log.d("Observe", "Observation started established")},
            {
                Handler(Looper.getMainLooper()).post {
                    Toast.makeText(this, it.item().toString(), Toast.LENGTH_LONG).show()
                }
                Log.d("item change", "item change"+it.item().toString())},
            {},
            {}
        )
    }
}
问题原因

直接通过控制台修改DynamoDB数据时,不会触发AWS AppSync的mutation事件——而你当前使用的Amplify.API.subscribe(基于AppSync订阅)和Amplify.DataStore.observe都依赖AppSync的变更通知机制,因此无法感知到这类直接修改数据库的操作。另外,DataStore默认的增量同步需要云端生成变更记录,才能同步到本地。

解决方案

方案1:配置DynamoDB Stream触发AppSync Mutation(推荐)

通过DynamoDB Stream捕获表的变更,再用Lambda函数调用AppSync的mutation接口,将变更同步到AppSync系统中,这样客户端订阅和DataStore就能正常收到通知:

  • 开启目标DynamoDB表的Stream(选择「新图像和旧图像」)
  • 创建Lambda函数,监听该Stream事件
  • 在Lambda中编写代码,调用AppSync的mutation接口来同步变更数据
  • 为Lambda配置调用AppSync的权限

方案2:使用DataStore的observeQuery实现实时同步

observeQuery会自动处理云端与本地的同步,定期拉取云端变更(默认30秒间隔,可配置),相比observe更适合监听包括控制台直接修改的场景。取消代码中注释的observeQuery部分,调整为以下配置:

val tag = "ObserveQuery"
val options = ObserveQueryOptions.builder()
    .predicate(Todo.NAME.beginsWith("Build")) // 保留你的过滤条件,或传null监听所有数据
    .syncMode(ObserveQueryOptions.SyncMode.REALTIME)
    .build()

Amplify.DataStore.observeQuery(
    Todo::class.java,
    options,
    { Log.d(tag, "监听已启动") },
    { snapshot ->
        // 接收初始数据和后续变更快照
        Log.d(tag, "当前数据总数: ${snapshot.items.size}")
        snapshot.items.forEach { todo ->
            Log.d(tag, "更新的Todo: ${todo.name}")
        }
        // 主线程更新UI
        Handler(Looper.getMainLooper()).post {
            Toast.makeText(this, "数据已更新", Toast.LENGTH_SHORT).show()
        }
    },
    { Log.e(tag, "监听失败", it) },
    { Log.d(tag, "监听结束") }
)

方案3:确保DataStore同步配置正确

  • 确认模型类上有@model注解,且生成时开启了同步(执行amplify add api时选择同步选项)
  • 监听DataStore同步状态,确认同步过程无异常:
Amplify.DataStore.sync(
    { Log.i("DataStore", "同步进度: 已同步${it.itemCount}条数据") },
    { Log.e("DataStore", "同步失败", it) },
    { Log.i("DataStore", "同步完成") }
)
代码优化建议
  • 避免重复监听:当前代码同时使用了2次Amplify.DataStore.observe和2次Amplify.API.subscribe,建议统一使用observeQuery简化逻辑
  • 管理订阅生命周期:API订阅需在Activity销毁时取消,避免内存泄漏:
private var updateSubscription: Cancelable? = null
private var createSubscription: Cancelable? = null

override fun onCreate(savedInstanceState: Bundle?) {
    // ... 其他代码
    updateSubscription = Amplify.API.subscribe(ModelSubscription.onUpdate(Todo::class.java), ...)
    createSubscription = Amplify.API.subscribe(ModelSubscription.onCreate(Todo::class.java), ...)
}

override fun onDestroy() {
    super.onDestroy()
    updateSubscription?.cancel()
    createSubscription?.cancel()
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 17:44:51