技术开发 频道

使用BMT消息驱动BEAN和SPRING进行息处理

重复发送侦测和预防
 
  在前面的部分中,我们看到,当非事务型地消费消息,并且要求once-and-only-once的服务质量时,应用程序需要引入特殊的机制来防止重复处理,这种重复处理在一个或多个组件在事务提交和消息通知两步之间发生故障这一罕见的场景下会发生。这种机制必须足够灵活以应付不同的部署场景;例如,如果应用程序部署在集群上,则重新发送消息的节点可能与最初发送此消息的节点不同。为了发现重复以防止消息被再次处理,应用程序应该能够通过某种方法识别消息已经被处理过了。
 
  在数据库实例中标识已经处理过的消息似乎比较合理,数据库实例(通常)被集群中所有节点共享。我们可以用JMSMessageID作为消息的独有标识符,根据JMS文档,JMSMessageID是一个字符串,它在历史库中是惟一标识消息的键。虽然惟一性的确切范围是由提供者定义的,但是它至少可以覆盖一个特定提供者安装中的所有消息,这里的安装指的是某种消息路由器的互连集合。为了安全起见,我们还将使用JMSTimestamp。JMSTimestamp字段包含了将消息转交给提供者以供发送的时间。这两个字段的值结合起来,就可以保证惟一地标识每条消息,即使它们来自不同的提供者。
 
  显然,尽可能地将重复侦测机制设计得快速且高效是非常有益的,所以我们会完全避免使用SQL选择操作,而是创建一个带有JMSMessageID和JMSTimestamp字段中定义的组合主键的表。然后当随后试图插入具有相同主键的值的时候,由于违反数据的完整性,就可以侦测到重复记录,这样就可以侦测并进行适当的处理。下面是一个Oracle数据库的DDL脚本:
CREATE TABLE message_log (message_id VARCHAR2(255) NOT NULL, message_ts INTEGER NOT NULL, CONSTRAINT pk_message_log PRIMARY KEY (message_id, message_ts));
实现非事务型消息消费者模型
 
  让我们尝试构建一套工具来简化非事务型消息的处理过程。由于涉及到许多横切功能,利用AOP方法会比较好。下面,我们将演示如何利用Spring Framework的AOP特性来实现具有once-and-only-once服务质量的非事务型消息处理所需的所有逻辑。
 
  Spring AOP Framework的强大功能在其他文章中已经有所论述。它使我们可以声明式地定义一组方面,每个方面为执行流提供特定的功能。产生的代码看起来非常结构化,并且每个部分可以独立测试,从而覆盖到特定于每个组件的所有用例。
 
  在Spring中使用AOP的首选方式是利用可以由Spring代理的组件接口。虽然JMS接口MessageListener和Message也可以用于此目的,但是定义一些更通用且与JMS接口解耦合的接口实际上会很有意义。例如,可以如下声明MessageProcessor接口:
public interface MessageProcessor { public Object process(MessageData messageData); }
  注意,process()方法带有一个MessageData参数,这个参数封装了来自原始消息的所有数据。可以实现Spring的MessageConverter,将Message实例转换成MessageData,并扩展AbstractJmsMessageDrivenBean来实现通用处理逻辑,即:
 
JMS消息到MessageData的转换
 
·从应用程序上下文获得处理器实例
·消息处理器的执行
·异常处理
 
以下是MessageDataDrivenBean的可能实现:
public abstract class MessageDataDrivenBean extends AbstractJmsMessageDrivenBean { private MessageConverter messageConverter = new MessageDataConverter(); private MessageProcessor messageProcessor; protected void onEjbCreate() { messageProcessor = ((MessageProcessor)   getBeanFactory().getBean(getName())); } public void onMessage(Message message) { try {    MessageData messageData = (MessageData) getMessageConverter().fromMessage(message);    this.messageProcessor.process(messageData); }   catch( MessageConversionException ex) {                    String msg = "Message conversion error;"+ex.getMessage();                    this.logger.error(msg, ex);                      throw new RuntimeException(msg, ex);                  }   catch( JMSException ex) {              String msg = "JMS error; "+ex.getMessage();              this.logger.error(msg, ex);              throw new RuntimeException(msg, ex);              } } protected MessageConverter getMessageConverter()   {    return this.messageConverter;   } protected abstract String getName(); }
  具体的子类只需要提供一个返回getName()方法的实现,当在Spring上下文中定义MessageProcessor接口时,该方法返回接口特定实现的bean名字。子类也可以提供自己的MessageConverter,使用不同的策略填充MessageData。以下是MessageProcessor的一个简单实现:
public class SimpleMessageProcessor implements MessageProcessor { private Log log = LogFactory.getLog(getClass()); public Object process( MessageData messageData) { log.info(messageData); return null; } }
  注意,虽然有可能为重复处理通知指定任意数据源,但是最好还是考虑一下当前的底层业务代码在做什么。例如,如果业务代码也使用了一个数据源,最好在通知中使用同一个——这样可能会减少事务中使用的不同XAResource数,而且,如果没有使用其他XAResource的话,J2EE服务器甚至可以使用单阶段提交优化,从而实际上使事务获得与本地相当的性能。
 
结束语
 
  正如您所看到的,我们成功实现了以非事务型方式处理消息消费并保持once-and-only-once服务质量的目标。我们相信在转向BMT以后,应用程序开发人员将获得更多对应用程序行为和事务处理的控制权,并同时享受明显的性能提升。这一转换的另一个作用是,如果只打算消费消息,那么您不必关心所使用的MOM是否具有对XA事务的健壮支持。
 
  解决方案可以使用Spring的AOP框架以一种非常结构化的方法实现,这使您可以利用一些相互独立的通知将一些非功能需求的单个方面组合起来。每个通知负责处理功能的一个特定方面。但是,将其组装在一个Spring应用程序上下文中,就可以管理消息处理的所有非功能性需求,包括事物处理、重复处理、被污染消息的处理,以及常见的跟踪、日志记录和性能监控。
0
相关文章