如何在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
相关产品推荐
相关产品推荐

