【发布时间】:2013-01-07 09:06:51
【问题描述】:
我想在我的 Java 程序中传递一条异步消息,所以第一步它应该持续监控数据库中某些表的变化。当有新的传入消息时,它应该显示它。只要应用程序正在运行,这应该是重复的过程。
我可以知道如何处理以下代码,其中包含轮询方法,它必须每 6 秒无限地调用自身,并且还应该在数据库中找到新的传入消息。
这里是sn-p的代码:
public class PollingSynchronizer implements Runnable {
private Collection<KPIMessage> incomingMessages;
private Connection dbConnection;
/**
* Constructor. Requires to provide a reference to the KA message queue
*
* @param incomingMessages reference to message queue
*
*/
public PollingSynchronizer(Collection<KpiMessage> incomingMessages, Connection dbConnection) {
super();
this.incomingMessages = incomingMessages;
this.dbConnection = dbConnection;
}
private int sequenceId;
public int getSequenceId() {
return sequenceId;
}
public void setSequenceId(int sequenceId) {
this.sequenceId = sequenceId;
}
@Override
/**
* The method which runs Polling action and record the time at which it is done
*
*/
public void run() {
try {
incomingMessages.addAll(fullPoll());
System.out.println("waiting 6 seconds");
//perform this operation in a loop
Thread.sleep(6000);
} catch (InterruptedException e) {
// TODO Auto-generated catch block
e.printStackTrace();
} catch (Exception e) {
// TODO Auto-generated catch block
e.printStackTrace();
}
Date currentDate = new Date();
SimpleDateFormat sdf = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss.SSSSSS");
// System.out.println(sdf.format(currentDate) + " " + msg);
}
/**
* Method which defines polling of the database and also count the number of Queries
* @return
* @throws Exception
*/
public List<KpiMessage> fullPoll() throws Exception {
// int sequenceID = 0;
Statement st = dbConnection.createStatement();
ResultSet rs = st.executeQuery("select * from msg_new_to_bde where ACTION = 804 order by SEQ DESC");
List<KpiMessage> pojoCol = new ArrayList<KpiMessage>();
while (rs.next()) {
KpiMessage filedClass = convertRecordsetToPojo(rs);
pojoCol.add(filedClass);
}
return pojoCol;
}
/**
* Converts a provided record-set to a {@link KpiMessage}.
*
* The following attributes are copied from record-set to pojo:
*
* <ul>
* <li>SEQ</li>
* <li>TABLENAME</li>
* <li>ENTRYTIME</li>
* <li>STATUS</li>
* </ul>
*
* @param rs
* the recordset to convert
* @return the converted pojo class object
* @throws SQLException
* if an sql error occurrs during processing of recordset
*/
private KpiMessage convertRecordsetToPojo(ResultSet rs) throws SQLException {
KpiMessage msg = new KpiMessage();
int sequence = rs.getInt("SEQ");
msg.setSequence(sequence);
int action = rs.getInt("ACTION");
msg.setAction(action);
String tablename = rs.getString("TABLENAME");
msg.setTableName(tablename);
Timestamp entrytime = rs.getTimestamp("ENTRYTIME");
Date entryTime = new Date(entrytime.getTime());
msg.setEntryTime(entryTime);
Timestamp processingtime = rs.getTimestamp("PROCESSINGTIME");
if (processingtime != null) {
Date processingTime = new Date(processingtime.getTime());
msg.setProcessingTime(processingTime);
}
String keyInfo1 = rs.getString("KEYINFO1");
msg.setKeyInfo1(keyInfo1);
String keyInfo2 = rs.getString("KEYINFO2");
msg.setKeyInfo2(keyInfo2);
return msg;
}
}
这里的序列 id 是表中的唯一 id,它随着新的传入消息的到达而不断增加。
P.S :“恳请:请给出给出负分的理由(大拇指向下)。这样我就可以清楚地解释我的问题了”
【问题讨论】:
-
我 - 我会使用 Quartz 创建一个重复的轮询任务,我会使用 JMS(更具体地说是 HornetQ)来处理消息传递部分。当已经有坚如磐石的轮子可用时,我不喜欢重新发明轮子。
-
投票?你可以通过触发器引发某种事件吗?
-
@thanks Gimby 我可以知道如何使用石英,这也是一个初学者级别的程序员,还如何使用异步调用这个轮询,就像一个人不断轮询和其他处理消息并更新它...... .
-
@MartinJames 我不知道,我只是卡在这里....
-
当您没有活动时,轮询是另一种方式。您通过阅读它的手册来学习使用石英,这在一篇文章中无法回答。
标签: java multithreading asynchronous parallel-processing polling