本文目录导读:

我来详细解释Java中分布式数据适配器模式的实现方式。
分布式数据适配器模式概述
在分布式系统中,适配器模式用于统一不同数据源、不同协议、不同格式的数据访问接口。
基础适配器实现
1 统一数据访问接口
public interface DataSourceAdapter {
<T> T query(String sql, Class<T> resultType);
<T> List<T> queryList(String sql, Class<T> resultType);
int execute(String sql);
<T> T executeWithTransaction(Callback<T> callback);
}
public interface Callback<T> {
T execute(DataSourceAdapter adapter);
}
2 多数据源适配器
// MySQL适配器
public class MySQLAdapter implements DataSourceAdapter {
private final DataSource dataSource;
public MySQLAdapter(DataSource dataSource) {
this.dataSource = dataSource;
}
@Override
public <T> T query(String sql, Class<T> resultType) {
// MySQL特定实现
try (Connection conn = dataSource.getConnection();
PreparedStatement stmt = conn.prepareStatement(sql)) {
ResultSet rs = stmt.executeQuery();
return mapResult(rs, resultType);
} catch (SQLException e) {
throw new DataAccessException("MySQL query failed", e);
}
}
@Override
public <T> List<T> queryList(String sql, Class<T> resultType) {
// 实现查询列表
return new ArrayList<>();
}
@Override
public int execute(String sql) {
// 执行更新操作
return 0;
}
@Override
public <T> T executeWithTransaction(Callback<T> callback) {
// 事务处理
return null;
}
}
// Redis适配器
public class RedisAdapter implements DataSourceAdapter {
private final JedisPool jedisPool;
public RedisAdapter(JedisPool jedisPool) {
this.jedisPool = jedisPool;
}
@Override
public <T> T query(String key, Class<T> resultType) {
try (Jedis jedis = jedisPool.getResource()) {
String value = jedis.get(key);
return deserialize(value, resultType);
}
}
@Override
public <T> List<T> queryList(String pattern, Class<T> resultType) {
try (Jedis jedis = jedisPool.getResource()) {
Set<String> keys = jedis.keys(pattern);
List<T> results = new ArrayList<>();
for (String key : keys) {
String value = jedis.get(key);
results.add(deserialize(value, resultType));
}
return results;
}
}
}
// MongoDB适配器
public class MongoDBAdapter implements DataSourceAdapter {
private final MongoCollection<Document> collection;
public MongoDBAdapter(MongoClient mongoClient, String dbName, String collectionName) {
MongoDatabase database = mongoClient.getDatabase(dbName);
this.collection = database.getCollection(collectionName);
}
@Override
public <T> T query(String filterJson, Class<T> resultType) {
Bson filter = Document.parse(filterJson);
Document doc = collection.find(filter).first();
return doc != null ? docToObject(doc, resultType) : null;
}
}
分布式缓存适配器
// 缓存适配器接口
public interface CacheAdapter {
<T> T get(String key, Class<T> type);
void put(String key, Object value, long ttl);
void evict(String key);
void clear();
}
// 多级缓存适配器
public class MultiLevelCacheAdapter implements CacheAdapter {
private final CacheAdapter localCache; // 本地缓存
private final CacheAdapter redisCache; // Redis缓存
private final CacheAdapter distributedCache; // 分布式缓存
public MultiLevelCacheAdapter() {
this.localCache = new LocalCacheAdapter();
this.redisCache = new RedisCacheAdapter();
this.distributedCache = new HazelcastCacheAdapter();
}
@Override
public <T> T get(String key, Class<T> type) {
// 1. 先查本地缓存
T value = localCache.get(key, type);
if (value != null) return value;
// 2. 再查Redis
value = redisCache.get(key, type);
if (value != null) {
localCache.put(key, value, 300); // 写入本地缓存
return value;
}
// 3. 最后查分布式缓存
value = distributedCache.get(key, type);
if (value != null) {
redisCache.put(key, value, 3600);
localCache.put(key, value, 300);
}
return value;
}
@Override
public void put(String key, Object value, long ttl) {
// 更新所有缓存层
localCache.put(key, value, ttl);
redisCache.put(key, value, ttl);
distributedCache.put(key, value, ttl);
}
}
数据库读写分离适配器
public class ReadWriteSeparationAdapter implements DataSourceAdapter {
private final DataSourceAdapter writeAdapter; // 写库
private final List<DataSourceAdapter> readAdapters; // 读库列表
private final LoadBalancer loadBalancer;
public ReadWriteSeparationAdapter(DataSourceAdapter writeAdapter,
List<DataSourceAdapter> readAdapters) {
this.writeAdapter = writeAdapter;
this.readAdapters = readAdapters;
this.loadBalancer = new RoundRobinLoadBalancer();
}
@Override
public <T> T query(String sql, Class<T> resultType) {
// 读操作路由到从库
DataSourceAdapter readAdapter = loadBalancer.select(readAdapters);
return readAdapter.query(sql, resultType);
}
@Override
public int execute(String sql) {
// 写操作路由到主库
return writeAdapter.execute(sql);
}
@Override
public <T> T executeWithTransaction(Callback<T> callback) {
// 事务操作使用写库
return writeAdapter.executeWithTransaction(callback);
}
// 负载均衡器
interface LoadBalancer {
<T> T select(List<T> adapters);
}
static class RoundRobinLoadBalancer implements LoadBalancer {
private AtomicInteger counter = new AtomicInteger(0);
@Override
public <T> T select(List<T> adapters) {
int index = Math.abs(counter.getAndIncrement()) % adapters.size();
return adapters.get(index);
}
}
}
分布式事务适配器
public class DistributedTransactionAdapter {
private final TransactionManager transactionManager;
private final Map<String, DataSourceAdapter> participants;
public DistributedTransactionAdapter(TransactionManager transactionManager) {
this.transactionManager = transactionManager;
this.participants = new ConcurrentHashMap<>();
}
public void registerParticipant(String name, DataSourceAdapter adapter) {
participants.put(name, adapter);
}
public <T> T executeGlobalTransaction(GlobalTransactionCallback<T> callback) {
// 开始全局事务
String transactionId = transactionManager.beginTransaction();
try {
T result = callback.execute(transactionId, participants);
transactionManager.commit(transactionId);
return result;
} catch (Exception e) {
transactionManager.rollback(transactionId);
throw new TransactionException("Global transaction failed", e);
}
}
@FunctionalInterface
interface GlobalTransactionCallback<T> {
T execute(String transactionId, Map<String, DataSourceAdapter> participants);
}
}
协议适配器
// REST API适配器
public class RestApiAdapter implements DataSourceAdapter {
private final RestTemplate restTemplate;
private final String baseUrl;
public RestApiAdapter(RestTemplate restTemplate, String baseUrl) {
this.restTemplate = restTemplate;
this.baseUrl = baseUrl;
}
@Override
public <T> T query(String endpoint, Class<T> resultType) {
String url = baseUrl + endpoint;
ResponseEntity<T> response = restTemplate.getForEntity(url, resultType);
return response.getBody();
}
@Override
public int execute(String endpoint) {
String url = baseUrl + endpoint;
ResponseEntity<Void> response = restTemplate.postForEntity(url, null, Void.class);
return response.getStatusCode().value();
}
}
// gRPC适配器
public class GrpcAdapter implements DataSourceAdapter {
private final ManagedChannel channel;
private final DataServiceGrpc.DataServiceBlockingStub stub;
public GrpcAdapter(String host, int port) {
this.channel = ManagedChannelBuilder.forAddress(host, port)
.usePlaintext()
.build();
this.stub = DataServiceGrpc.newBlockingStub(channel);
}
@Override
public <T> T query(String request, Class<T> resultType) {
DataRequest dataRequest = DataRequest.newBuilder()
.setQuery(request)
.build();
DataResponse response = stub.query(dataRequest);
return deserializeResponse(response, resultType);
}
}
数据格式适配器
// JSON/XML/YAML格式适配
public interface DataFormatAdapter {
<T> T deserialize(String data, Class<T> type);
String serialize(Object data);
}
public class JsonFormatAdapter implements DataFormatAdapter {
private final ObjectMapper objectMapper = new ObjectMapper();
@Override
public <T> T deserialize(String data, Class<T> type) {
try {
return objectMapper.readValue(data, type);
} catch (JsonProcessingException e) {
throw new DataFormatException("JSON deserialization failed", e);
}
}
@Override
public String serialize(Object data) {
try {
return objectMapper.writeValueAsString(data);
} catch (JsonProcessingException e) {
throw new DataFormatException("JSON serialization failed", e);
}
}
}
public class ProtobufFormatAdapter implements DataFormatAdapter {
@Override
public <T> T deserialize(String data, Class<T> type) {
// Protobuf反序列化实现
return null;
}
@Override
public String serialize(Object data) {
// Protobuf序列化实现
return null;
}
}
使用示例
public class DataAdapterExample {
public static void main(String[] args) {
// 创建数据源适配器
DataSourceAdapter mysqlAdapter = new MySQLAdapter(createDataSource());
DataSourceAdapter redisAdapter = new RedisAdapter(createJedisPool());
DataSourceAdapter mongoAdapter = new MongoDBAdapter(createMongoClient(),
"database", "collection");
// 创建多级缓存适配器
CacheAdapter multiLevelCache = new MultiLevelCacheAdapter();
// 创建读写分离适配器
List<DataSourceAdapter> readAdapters = Arrays.asList(
new MySQLAdapter(createReadDataSource1()),
new MySQLAdapter(createReadDataSource2())
);
DataSourceAdapter readWriteAdapter = new ReadWriteSeparationAdapter(
mysqlAdapter, readAdapters
);
// 使用适配器
User user = readWriteAdapter.query("SELECT * FROM users WHERE id = 1", User.class);
// 使用缓存
User cachedUser = multiLevelCache.get("user:1", User.class);
if (cachedUser == null) {
cachedUser = mysqlAdapter.query("SELECT * FROM users WHERE id = 1", User.class);
multiLevelCache.put("user:1", cachedUser, 3600);
}
}
}
最佳实践
- 接口统一:定义通用的适配器接口,屏蔽底层差异
- 依赖注入:使用依赖注入管理适配器实例
- 连接池管理:统一管理数据库连接池
- 监控和日志:添加性能监控和日志记录
- 故障转移:实现适配器的自动故障转移机制
- 配置中心:使用配置中心管理适配器配置
这样的适配器模式可以让分布式系统灵活应对不同的数据源和协议,同时保持代码的一致性和可维护性。