GH-3578: Fix JdbcMessageStore.getMessageGroup() (#3579)
Fixes https://github.com/spring-projects/spring-integration/issues/3578 The `JdbcTemplate.queryForMap()` extract values for columns to the closer target driver types. For example H2 and Derby return `Long` for `BIGINT`. Oracle for its `NUMBER(19,0)` returns `BigInteger`. This makes the code in the `JdbcMessageStore.getMessageGroup()` not platform independent. * Fix `JdbcMessageStore` to map `ResultSet` to the `MessageGroupMetadata` directly. Mostly reinstating the previous behavior * For that reason expose a default ctor for `MessageGroupMetadata` and extract some setters to make code in the `JdbcMessageStore` more cleaner. * This opens for us a possibility to implement a `MessageGroupStore.getGroupMetadata(groupId)` for `JdbcMessageStore` * Fix deprecation for `Flux.limitRequest()`
This commit is contained in:
@@ -391,10 +391,10 @@ public abstract class AbstractPollingEndpoint extends AbstractEndpoint implement
|
||||
fluxSink.complete();
|
||||
}
|
||||
})
|
||||
.limitRequest(
|
||||
this.maxMessagesPerPoll < 0
|
||||
.take(this.maxMessagesPerPoll < 0
|
||||
? Long.MAX_VALUE
|
||||
: this.maxMessagesPerPoll);
|
||||
: this.maxMessagesPerPoll,
|
||||
true);
|
||||
}
|
||||
})
|
||||
.subscribeOn(Schedulers.fromExecutor(this.taskExecutor))
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2019 the original author or authors.
|
||||
* Copyright 2002-2021 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
@@ -40,7 +40,7 @@ public class MessageGroupMetadata implements Serializable {
|
||||
|
||||
private static final long serialVersionUID = 1L;
|
||||
|
||||
private List<UUID> messageIds = new LinkedList<>();
|
||||
private final List<UUID> messageIds = new LinkedList<>();
|
||||
|
||||
private long timestamp;
|
||||
|
||||
@@ -52,8 +52,7 @@ public class MessageGroupMetadata implements Serializable {
|
||||
|
||||
private volatile String condition;
|
||||
|
||||
private MessageGroupMetadata() {
|
||||
//For Jackson deserialization
|
||||
public MessageGroupMetadata() {
|
||||
}
|
||||
|
||||
public MessageGroupMetadata(MessageGroup messageGroup) {
|
||||
@@ -79,7 +78,7 @@ public class MessageGroupMetadata implements Serializable {
|
||||
return !this.messageIds.contains(messageId) && this.messageIds.add(messageId);
|
||||
}
|
||||
|
||||
void setLastModified(long lastModified) {
|
||||
public void setLastModified(long lastModified) {
|
||||
this.lastModified = lastModified;
|
||||
}
|
||||
|
||||
@@ -107,7 +106,7 @@ public class MessageGroupMetadata implements Serializable {
|
||||
return new LinkedList<UUID>(this.messageIds);
|
||||
}
|
||||
|
||||
void complete() {
|
||||
public void complete() {
|
||||
this.complete = true;
|
||||
}
|
||||
|
||||
@@ -123,11 +122,15 @@ public class MessageGroupMetadata implements Serializable {
|
||||
return this.timestamp;
|
||||
}
|
||||
|
||||
public void setTimestamp(long timestamp) {
|
||||
this.timestamp = timestamp;
|
||||
}
|
||||
|
||||
public int getLastReleasedMessageSequenceNumber() {
|
||||
return this.lastReleasedMessageSequenceNumber;
|
||||
}
|
||||
|
||||
void setLastReleasedMessageSequenceNumber(int lastReleasedMessageSequenceNumber) {
|
||||
public void setLastReleasedMessageSequenceNumber(int lastReleasedMessageSequenceNumber) {
|
||||
this.lastReleasedMessageSequenceNumber = lastReleasedMessageSequenceNumber;
|
||||
}
|
||||
|
||||
|
||||
@@ -21,7 +21,6 @@ import java.sql.SQLException;
|
||||
import java.sql.Timestamp;
|
||||
import java.util.Arrays;
|
||||
import java.util.Collection;
|
||||
import java.util.Collections;
|
||||
import java.util.HashMap;
|
||||
import java.util.Iterator;
|
||||
import java.util.List;
|
||||
@@ -38,6 +37,7 @@ import org.springframework.dao.DuplicateKeyException;
|
||||
import org.springframework.dao.IncorrectResultSizeDataAccessException;
|
||||
import org.springframework.integration.store.AbstractMessageGroupStore;
|
||||
import org.springframework.integration.store.MessageGroup;
|
||||
import org.springframework.integration.store.MessageGroupMetadata;
|
||||
import org.springframework.integration.store.MessageMetadata;
|
||||
import org.springframework.integration.store.MessageStore;
|
||||
import org.springframework.integration.store.SimpleMessageGroup;
|
||||
@@ -334,11 +334,13 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa
|
||||
@Override
|
||||
public void addMessagesToGroup(Object groupId, Message<?>... messages) {
|
||||
String groupKey = getKey(groupId);
|
||||
Map<String, Object> groupInfo = getGroupMetadata(groupKey);
|
||||
MessageGroupMetadata groupMetadata = getGroupMetadata(groupKey);
|
||||
|
||||
Timestamp updatedDate = new Timestamp(System.currentTimeMillis());
|
||||
boolean groupNotExist = groupInfo.isEmpty();
|
||||
Timestamp createdDate = groupNotExist ? updatedDate : (Timestamp) groupInfo.get("CREATED_DATE");
|
||||
boolean groupNotExist = groupMetadata == null;
|
||||
Timestamp createdDate =
|
||||
groupNotExist
|
||||
? new Timestamp(System.currentTimeMillis())
|
||||
: new Timestamp(groupMetadata.getTimestamp());
|
||||
|
||||
for (Message<?> message : messages) {
|
||||
addMessage(message);
|
||||
@@ -397,28 +399,39 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa
|
||||
|
||||
@Override
|
||||
public MessageGroup getMessageGroup(Object groupId) {
|
||||
String key = getKey(groupId);
|
||||
Map<String, Object> groupInfo = getGroupMetadata(key);
|
||||
|
||||
if (groupInfo.isEmpty()) {
|
||||
MessageGroupMetadata groupMetadata = getGroupMetadata(groupId);
|
||||
if (groupMetadata != null) {
|
||||
MessageGroup messageGroup =
|
||||
getMessageGroupFactory()
|
||||
.create(this, groupId, groupMetadata.getTimestamp(), groupMetadata.isComplete());
|
||||
messageGroup.setLastModified(groupMetadata.getLastModified());
|
||||
messageGroup.setLastReleasedMessageSequenceNumber(groupMetadata.getLastReleasedMessageSequenceNumber());
|
||||
messageGroup.setCondition(groupMetadata.getCondition());
|
||||
return messageGroup;
|
||||
}
|
||||
else {
|
||||
return new SimpleMessageGroup(groupId);
|
||||
}
|
||||
|
||||
MessageGroup messageGroup = getMessageGroupFactory()
|
||||
.create(this, groupId, ((Timestamp) groupInfo.get("CREATED_DATE")).getTime(),
|
||||
((Long) groupInfo.get("COMPLETE")) > 0);
|
||||
messageGroup.setLastModified(((Timestamp) groupInfo.get("UPDATED_DATE")).getTime());
|
||||
messageGroup.setLastReleasedMessageSequenceNumber(((Long) groupInfo.get("LAST_RELEASED_SEQUENCE")).intValue());
|
||||
messageGroup.setCondition((String) groupInfo.get("CONDITION"));
|
||||
return messageGroup;
|
||||
}
|
||||
|
||||
private Map<String, Object> getGroupMetadata(String groupKey) {
|
||||
@Override
|
||||
public MessageGroupMetadata getGroupMetadata(Object groupId) {
|
||||
String key = getKey(groupId);
|
||||
try {
|
||||
return this.jdbcTemplate.queryForMap(getQuery(Query.GET_GROUP_INFO), groupKey, this.region);
|
||||
return this.jdbcTemplate.queryForObject(getQuery(Query.GET_GROUP_INFO), (rs, rowNum) -> {
|
||||
MessageGroupMetadata groupMetadata = new MessageGroupMetadata();
|
||||
if (rs.getInt("COMPLETE") > 0) {
|
||||
groupMetadata.complete();
|
||||
}
|
||||
groupMetadata.setTimestamp(rs.getTimestamp("CREATED_DATE").getTime());
|
||||
groupMetadata.setLastModified(rs.getTimestamp("UPDATED_DATE").getTime());
|
||||
groupMetadata.setLastReleasedMessageSequenceNumber(rs.getInt("LAST_RELEASED_SEQUENCE"));
|
||||
groupMetadata.setCondition(rs.getString("CONDITION"));
|
||||
return groupMetadata;
|
||||
}, key, this.region);
|
||||
}
|
||||
catch (IncorrectResultSizeDataAccessException ex) {
|
||||
return Collections.emptyMap();
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user