【发布时间】:2016-10-11 08:34:05
【问题描述】:
我有一个聊天应用程序,它使用 socket.io 与 nodejs 服务器(由我编写)通信。多人可以在一对一的基础上互相聊天。该应用程序的用户界面类似于 WhatsApp。有不同的Chat Threads,其中Users可以交换Chat Messages。这些需要存储在 SQLite 数据库中。
为确保应用即使在关闭时也能正常工作,我已在服务中启动了 socket.io 连接。
当有新消息到来时,顺序是这样的
- 服务解析响应并将其写入 SQLite db
- 如果应用关闭,则会创建通知并向用户显示
- 如果应用已打开,并且用户在发送新消息的同一聊天线程中,则必须更新
recyclerView。 - 基本上感觉和whatsapp一样,除了我用的是socket.io而不是XMPP
我想使用RxJava 对SQLite 执行CRUD 操作。
SQLiteHelper.Class
public List<ChatMessage> getChatMessagesForThread(String threadName){
List<ChatMessage>list = new ArrayList<>();
String query = "Some Query Here";
SQLiteDatabase db = this.getReadableDatabase();
Cursor c = db.rawQuery(query, null);
if (c.moveToFirst()){
do {
ChatMessage message = new ChatMessage();
message.setMessage(c.getString((c.getColumnIndex(CHAT_MESSAGES_KEY_MESSAGE))));
message.setChatThread(c.getInt(c.getColumnIndex(CHAT_MESSAGES_KEY_CHAT_THREAD)));
message.setUser(c.getString(c.getColumnIndex(CHAT_MESSAGES_KEY_USER)));
list.add(message);
} while (c.moveToNext());
}else {
//No such thread exists. Returns null
}
c.close();
return list;
}
Rx Java 相关函数
public static <T> Observable<T> makeObservable(final Callable<T> func) {
return Observable.create(
new Observable.OnSubscribe<T>() {
@Override
public void call(Subscriber<? super T> subscriber) {
try {
T observed = func.call();
if (observed != null) { // to make defaultIfEmpty work
subscriber.onNext(observed);
}
//TODO: DECIDE IF THIS STAYS OR NOT !!!
//subscriber.onCompleted();
} catch (Exception ex) {
subscriber.onError(ex);
}
}
}).subscribeOn(Schedulers.io());
}
@SuppressWarnings("unchecked")
private Callable<List<ChatMessage>> getData(String threadName) {
return new Callable() {
public List<ChatMessage> call() {
return getChatMessagesForThread(String threadName);
}
};
}
public Observable<List<ChatMessage>> getDataObservable(String threadName) {
return makeObservable(getData(String threadName));
}
在ChatActivity.class
databaseHelper.getDataObservable(String currentThreadName)
.subscribe(new Action1<List<ChatMessage>>() {
@Override
public void call(List<ChatMessage> chatMessages) {
for (ChatMessage message : chatMessages){
//Add to list and update recyclerView
mAdapter.notifyDatasetChanged();
}
}
});
问题是,上面的 Rxjava 调用只执行一次,而不是每次完成新的写入操作。
此函数在 SQLiteHelper 类中,从我的服务调用
public long writeMessageToDB(ChatMessage message){
ContentValues values = new ContentValues();
values.put(CHAT_MESSAGES_KEY_MESSAGE, message.getMessage());
values.put(CHAT_MESSAGES_KEY_CHAT_THREAD, message.getChatThread());
values.put(CHAT_MESSAGES_KEY_USER, message.getUser());
SQLiteDatabase db = this.getWritableDatabase();
return db.insert(CHAT_MESSAGES_TABLE_NAME, null, values);
}
我已阅读 Square 的 SQL Delite,但想了解如何使上述内容正常工作。请帮忙!
【问题讨论】:
标签: android sqlite android-contentprovider rx-java reactive-programming