最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

Java基于SpringBoot和tk.mybatis實現(xiàn)事務(wù)讀寫分離代碼實例

 更新時間:2023年10月07日 10:00:31   作者:sunct  
這篇文章主要介紹了Java基于SpringBoot和tk.mybatis實現(xiàn)事務(wù)讀寫分離代碼實例,讀寫分離,基本的原理是讓主數(shù)據(jù)庫處理事務(wù)性增、改、刪操作,而從數(shù)據(jù)庫處理SELECT查詢操作,數(shù)據(jù)庫復(fù)制被用來把事務(wù)性操作導(dǎo)致的變更同步到集群中的從數(shù)據(jù)庫,需要的朋友可以參考下

什么是讀寫分離?

讀寫分離,基本的原理是讓主數(shù)據(jù)庫處理事務(wù)性增、改、刪操作( INSERT、UPDATE、 DELETE) ,而從數(shù)據(jù)庫處理SELECT查詢操作。

數(shù)據(jù)庫復(fù)制被用來把事務(wù)性操作導(dǎo)致的變更同步到集群中的從數(shù)據(jù)庫。

為什么要讀寫分離呢?

  • 因為數(shù)據(jù)庫的“寫”(寫10000條數(shù)據(jù)可能要3分鐘)操作是比較耗時的。
  • 但是數(shù)據(jù)庫的“讀”(讀10000條數(shù)據(jù)可能只要5秒鐘)
  • 所以讀寫分離,解決的是,數(shù)據(jù)庫的寫入,影響了查詢的效率。

源碼

先定義數(shù)據(jù)源讀寫類型

/**
 * 數(shù)據(jù)源類型
 *
 * @author sunchangtan
 */
public enum DataSourceType {
    WRITE, READ
}

定義數(shù)據(jù)庫連接的Holder

import java.sql.Connection;
import java.util.HashMap;
import java.util.Map;
import org.springframework.core.NamedThreadLocal;
/**
 * 數(shù)據(jù)庫連接的Holder
 *
 * @author sunchangtan
 */
public class ConnectionHolder {
    /**
     * 當(dāng)前數(shù)據(jù)庫鏈接是讀還是寫
     */
    public final static ThreadLocal<DataSourceType> CURRENT_CONNECTION = new NamedThreadLocal<DataSourceType>("routingdatasource's key") {
        protected DataSourceType initialValue() {
            return DataSourceType.WRITE;
        }
    };
    /**
     * 當(dāng)前線程所有數(shù)據(jù)庫鏈接
     */
    public final static ThreadLocal<Map<DataSourceType, Connection>> CONNECTION_CONTEXT = new NamedThreadLocal<Map<DataSourceType, Connection>>("connection map") {
        protected Map<DataSourceType, Connection> initialValue() {
            return new HashMap<>();
        }
    };
    /**
     * 強制寫數(shù)據(jù)源
     */
    public final static ThreadLocal<Boolean> FORCE_WRITE = new NamedThreadLocal<Boolean>("FORCE_WRITE");
}

定義數(shù)據(jù)源的Holder

import org.springframework.core.NamedThreadLocal;
/**
 * 數(shù)據(jù)源的Holder
 *
 * @author sunchangtan
 */
public class DataSourceHolder {
	/**
	 * 當(dāng)前數(shù)據(jù)組
	 */
	public final static ThreadLocal<DataSourceType> CURRENT_DATASOURCE = new NamedThreadLocal<>("routingdatasource's key");
	static {
		setCurrentDataSource(DataSourceType.WRITE);
	}
	public static void setCurrentDataSource(DataSourceType dataSourceType){
		CURRENT_DATASOURCE.set(dataSourceType);
	}
	public static DataSourceType getCurrentDataSource(){
		return CURRENT_DATASOURCE.get();
	}
	public static void clearDataSource() {
		CURRENT_DATASOURCE.remove();
	}
}

定義數(shù)據(jù)源代理類,處理讀寫數(shù)據(jù)庫的路由

import java.io.PrintWriter;
import java.lang.reflect.InvocationHandler;
import java.lang.reflect.InvocationTargetException;
import java.lang.reflect.Method;
import java.lang.reflect.Proxy;
import java.sql.Connection;
import java.sql.SQLException;
import java.util.Map;
import java.util.logging.Logger;
import javax.sql.DataSource;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.springframework.core.Constants;
import org.springframework.jdbc.datasource.ConnectionProxy;
/**
 * 數(shù)據(jù)庫代理,具體數(shù)據(jù)源由DataSourceRouter提供
 * 
 * @author sunchangtan
 */
public class DataSourceProxy implements DataSource {
	private static final Constants constants = new Constants(Connection.class);
	private static final Log logger = LogFactory.getLog(DataSourceProxy.class);
	private Boolean defaultAutoCommit = Boolean.TRUE;
	private Integer defaultTransactionIsolation = 2;
	private DataSourceRouter dataSourceRouter;
	public DataSourceProxy(DataSourceRouter dataSourceRouter) {
		this.dataSourceRouter = dataSourceRouter;
	}
	public void setDefaultAutoCommit(boolean defaultAutoCommit) {
		this.defaultAutoCommit = defaultAutoCommit;
	}
	public void setDefaultTransactionIsolation(int defaultTransactionIsolation) {
		this.defaultTransactionIsolation = defaultTransactionIsolation;
	}
	public void setDefaultTransactionIsolationName(String constantName) {
		setDefaultTransactionIsolation(constants.asNumber(constantName).intValue());
	}
	protected Boolean defaultAutoCommit() {
		return this.defaultAutoCommit;
	}
	protected Integer defaultTransactionIsolation() {
		return this.defaultTransactionIsolation;
	}
	@Override
	public Connection getConnection() throws SQLException {
		return (Connection) Proxy.newProxyInstance(ConnectionProxy.class.getClassLoader(),
				new Class<?>[] { ConnectionProxy.class }, new LazyConnectionInvocationHandler());
	}
	@Override
	public Connection getConnection(String username, String password) throws SQLException {
		return (Connection) Proxy.newProxyInstance(ConnectionProxy.class.getClassLoader(),
				new Class<?>[] { ConnectionProxy.class }, new LazyConnectionInvocationHandler(username, password));
	}
	private class LazyConnectionInvocationHandler implements InvocationHandler {
		private String username;
		private String password;
		private Boolean readOnly = Boolean.FALSE;
		private Integer transactionIsolation;
		private Boolean autoCommit;
		private boolean closed = false;
		public LazyConnectionInvocationHandler() {
			this.autoCommit = defaultAutoCommit();
			this.transactionIsolation = defaultTransactionIsolation();
		}
		public LazyConnectionInvocationHandler(String username, String password) {
			this();
			this.username = username;
			this.password = password;
		}
		@Override
		public Object invoke(Object proxy, Method method, Object[] args) throws Throwable {
			// Invocation on ConnectionProxy interface coming in...
			if (method.getName().equals("setTransactionIsolation") && args != null && (Integer) args[0] == Connection.TRANSACTION_SERIALIZABLE) {
				 args[0] = defaultTransactionIsolation();
				ConnectionHolder.FORCE_WRITE.set(Boolean.TRUE);
			}
			if (method.getName().equals("equals")) {
				// We must avoid fetching a target Connection for "equals".
				// Only consider equal when proxies are identical.
				return (proxy == args[0]);
			} else if (method.getName().equals("hashCode")) {
				// We must avoid fetching a target Connection for "hashCode",
				// and we must return the same hash code even when the target
				// Connection has been fetched: use hashCode of Connection
				// proxy.
				return System.identityHashCode(proxy);
			} else if (method.getName().equals("unwrap")) {
				if (((Class<?>) args[0]).isInstance(proxy)) {
					return proxy;
				}
			} else if (method.getName().equals("isWrapperFor")) {
				if (((Class<?>) args[0]).isInstance(proxy)) {
					return true;
				}
			} else if (method.getName().equals("getTargetConnection")) {
				// Handle getTargetConnection method: return underlying
				// connection.
				return getTargetConnection(method);
			}
			if (!hasTargetConnection()) {
				// No physical target Connection kept yet ->
				// resolve transaction demarcation methods without fetching
				// a physical JDBC Connection until absolutely necessary.
				if (method.getName().equals("toString")) {
					return "Lazy Connection proxy for target DataSource [" + dataSourceRouter.getTargetDataSource() + "]";
				} else if (method.getName().equals("getMetaData")) {
					return dataSourceRouter.getTargetDataSource().getConnection().getMetaData();
				} else if (method.getName().equals("isReadOnly")) {
					return this.readOnly;
				} else if (method.getName().equals("setReadOnly")) {
					this.readOnly = (Boolean) args[0];
					return null;
				} else if (method.getName().equals("getTransactionIsolation")) {
					if (this.transactionIsolation != null) {
						return this.transactionIsolation;
					}
					// Else fetch actual Connection and check there,
					// because we didn't have a default specified.
				} else if (method.getName().equals("setTransactionIsolation")) {
					this.transactionIsolation = (Integer) args[0];
					return null;
				} else if (method.getName().equals("getAutoCommit")) {
					if (this.autoCommit != null) {
						return this.autoCommit;
					}
					// Else fetch actual Connection and check there,
					// because we didn't have a default specified.
				} else if (method.getName().equals("setAutoCommit")) {
					this.autoCommit = (Boolean) args[0];
					return null;
				} else if (method.getName().equals("commit")) {
					// Ignore: no statements created yet.
					return null;
				} else if (method.getName().equals("rollback")) {
					// Ignore: no statements created yet.
					return null;
				} else if (method.getName().equals("getWarnings")) {
					return null;
				} else if (method.getName().equals("clearWarnings")) {
					return null;
				} else if (method.getName().equals("close")) {
					// Ignore: no target connection yet.
					this.closed = true;
					return null;
				} else if (method.getName().equals("isClosed")) {
					return this.closed;
				} else if (this.closed) {
					// Connection proxy closed, without ever having fetched a
					// physical JDBC Connection: throw corresponding
					// SQLException.
					throw new SQLException("Illegal operation: connection is closed");
				}
			} else {
				if (method.getName().equals("commit")) {
					Map<DataSourceType, Connection> connectionMap = ConnectionHolder.CONNECTION_CONTEXT.get();
					Connection writeCon = connectionMap.get(DataSourceType.WRITE);
					if (writeCon != null) {
						writeCon.commit();
					}
					return null;
				}
				if (method.getName().equals("rollback")) {
					Map<DataSourceType, Connection> connectionMap = ConnectionHolder.CONNECTION_CONTEXT.get();
					Connection writeCon = connectionMap.get(DataSourceType.WRITE);
					if (writeCon != null) {
						writeCon.rollback();
					}
					return null;
				}
				if (method.getName().equals("close")) {
		            ConnectionHolder.FORCE_WRITE.set(Boolean.FALSE);
					Map<DataSourceType, Connection> connectionMap = ConnectionHolder.CONNECTION_CONTEXT.get();
					Connection readCon = connectionMap.remove(DataSourceType.READ);
					if (readCon != null) {
					    readCon.close();
                    }
					Connection writeCon = connectionMap.remove(DataSourceType.WRITE);
					if (writeCon != null) {
						writeCon.close();
					}
					this.closed = true;
					return null;
				}
			}
			// Target Connection already fetched,
			// or target Connection necessary for current operation ->
			// invoke method on target connection.
			try {
			    return method.invoke(
	                     ConnectionHolder.CONNECTION_CONTEXT.get().get(ConnectionHolder.CURRENT_CONNECTION.get()), args);
			} catch (InvocationTargetException ex) {
				throw ex.getTargetException();
			}
		}
		/**
		 * Return whether the proxy currently holds a target Connection.
		 */
		private boolean hasTargetConnection() {
			return (ConnectionHolder.CONNECTION_CONTEXT.get() != null
					&& ConnectionHolder.CONNECTION_CONTEXT.get().get(ConnectionHolder.CURRENT_CONNECTION.get()) != null);
		}
		/**
		 * Return the target Connection, fetching it and initializing it if
		 * necessary.
		 */
		private Connection getTargetConnection(Method operation) throws SQLException {
			// No target Connection held -> fetch one.
			if (logger.isDebugEnabled()) {
				logger.debug("Connecting to database for operation '" + operation.getName() + "'");
			}
			// Fetch physical Connection from DataSource.
			Connection target = (this.username != null)
					? dataSourceRouter.getTargetDataSource().getConnection(this.username, this.password)
					: dataSourceRouter.getTargetDataSource().getConnection();
			// Apply kept transaction settings, if any.
			if (this.readOnly) {
				try {
					target.setReadOnly(this.readOnly);
				} catch (Exception ex) {
					// "read-only not supported" -> ignore, it's just a hint
					// anyway
					logger.debug("Could not set JDBC Connection read-only", ex);
				}
			}
			if (this.transactionIsolation != null && !this.transactionIsolation.equals(defaultTransactionIsolation())) {
				target.setTransactionIsolation(this.transactionIsolation);
			}
			if (DataSourceType.READ == ConnectionHolder.CURRENT_CONNECTION.get()) {
				try {
					target.setAutoCommit(true);
				} catch (SQLException e) {
					e.printStackTrace();
				}
			}
			if (this.autoCommit != null && this.autoCommit != target.getAutoCommit()) {
				if (DataSourceType.WRITE == ConnectionHolder.CURRENT_CONNECTION.get()) {
					target.setAutoCommit(this.autoCommit);
				}
			}
			return target;
		}
	}
	@Override
	public PrintWriter getLogWriter() throws SQLException {
		return dataSourceRouter.getTargetDataSource().getLogWriter();
	}
	@Override
	public void setLogWriter(PrintWriter out) throws SQLException {
		dataSourceRouter.getTargetDataSource().setLogWriter(out);
	}
	@Override
	public int getLoginTimeout() throws SQLException {
		return dataSourceRouter.getTargetDataSource().getLoginTimeout();
	}
	@Override
	public void setLoginTimeout(int seconds) throws SQLException {
		dataSourceRouter.getTargetDataSource().setLoginTimeout(seconds);
	}
	// ---------------------------------------------------------------------
	// Implementation of JDBC 4.0's Wrapper interface
	// ---------------------------------------------------------------------
	@Override
	@SuppressWarnings("unchecked")
	public <T> T unwrap(Class<T> iface) throws SQLException {
		if (iface.isInstance(this)) {
			return (T) this;
		}
		return dataSourceRouter.getTargetDataSource().unwrap(iface);
	}
	@Override
	public boolean isWrapperFor(Class<?> iface) throws SQLException {
		return (iface.isInstance(this) || dataSourceRouter.getTargetDataSource().isWrapperFor(iface));
	}
	// ---------------------------------------------------------------------
	// Implementation of JDBC 4.1's getParentLogger method
	// ---------------------------------------------------------------------
	@Override
	public Logger getParentLogger() {
		return Logger.getLogger(Logger.GLOBAL_LOGGER_NAME);
	}
}

定義路由接口

/**
 * 數(shù)據(jù)庫路由
 *
 */
public interface DataSourceRouter {
	/**
	 * 根據(jù)自己的需要,實現(xiàn)數(shù)據(jù)庫路由,可以是讀寫分離的數(shù)據(jù)源,或者是分表后的數(shù)據(jù)源
	 * @return
	 */
	public DataSource getTargetDataSource();
}

實現(xiàn)讀庫路由的基類

import org.springframework.beans.factory.InitializingBean;
import org.springframework.jdbc.datasource.lookup.DataSourceLookup;
import org.springframework.jdbc.datasource.lookup.JndiDataSourceLookup;
import javax.sql.DataSource;
import java.util.ArrayList;
import java.util.List;
/**
 *  讀寫數(shù)據(jù)源路由的基類
 *
 * @author sunchangtan
 */
public abstract class AbstractMasterSlaverDataSourceRouter implements DataSourceRouter, InitializingBean {
	// 配置文件中配置的read-only datasoure
	// 可以為真實的datasource,也可以jndi的那種
	private List<Object> readDataSources;
	private Object writeDataSource;
	private DataSourceLookup dataSourceLookup = new JndiDataSourceLookup();
	private List<DataSource> resolvedReadDataSources;
	private DataSource resolvedWriteDataSource;
	// read-only data source的數(shù)量,做負(fù)載均衡的時候需要
	private int readDsSize;
	public List<DataSource> getResolvedReadDataSources() {
		return resolvedReadDataSources;
	}
	public int getReadDsSize() {
		return readDsSize;
	}
	public void setReadDataSoures(List readDataSoures) {
		this.readDataSources = readDataSoures;
	}
	public void setWriteDataSource(Object writeDataSource) {
		this.writeDataSource = writeDataSource;
	}
	public void setDataSourceLookup(DataSourceLookup dataSourceLookup) {
		this.dataSourceLookup = (dataSourceLookup != null ? dataSourceLookup : new JndiDataSourceLookup());
	}
	@Override
	public void afterPropertiesSet() {
		if (writeDataSource == null) {
			throw new IllegalArgumentException("Property 'writeDataSource' is required");
		}
		this.resolvedWriteDataSource = resolveSpecifiedDataSource(writeDataSource);
		if (this.readDataSources == null || this.readDataSources.size() ==0) {
			throw new IllegalArgumentException("Property 'resolvedReadDataSources' is required");
		}
		resolvedReadDataSources = new ArrayList<DataSource>(readDataSources.size());
		for (Object item : readDataSources) {
			resolvedReadDataSources.add(resolveSpecifiedDataSource(item));
		}
		readDsSize = readDataSources.size();
	}
	protected DataSource resolveSpecifiedDataSource(Object dataSource) throws IllegalArgumentException {
		if (dataSource instanceof DataSource) {
			return (DataSource) dataSource;
		}
		else if (dataSource instanceof String) {
			return this.dataSourceLookup.getDataSource((String) dataSource);
		}
		else {
			throw new IllegalArgumentException(
					"Illegal data source value - only [javax.sql.DataSource] and String supported: " + dataSource);
		}
	}
	@Override
	public DataSource getTargetDataSource() {
		if (DataSourceType.WRITE.equals(ConnectionHolder.CURRENT_CONNECTION.get())) {
			return resolvedWriteDataSource;
		} else {
			return loadBalance();
		}
	}
	protected abstract DataSource loadBalance();
}

實現(xiàn)簡單的輪詢路由,其他路由方式,大家可以自行實現(xiàn)

import javax.sql.DataSource;
import java.util.concurrent.atomic.AtomicInteger;
/**
 * 
 * 簡單實現(xiàn)讀數(shù)據(jù)源負(fù)載均衡
 *
 */
public class RoundRobinMasterSlaverDataSourceRouter extends AbstractMasterSlaverDataSourceRouter {
	private AtomicInteger count = new AtomicInteger(0);
	@Override
	protected DataSource loadBalance() {
		int index = Math.abs(count.incrementAndGet()) % getReadDsSize();
		return getResolvedReadDataSources().get(index);
	}
}

處理“只讀事務(wù)到讀庫,讀寫事務(wù)到寫庫”的事務(wù)

import org.springframework.jdbc.datasource.DataSourceTransactionManager;
import org.springframework.transaction.TransactionDefinition;
import javax.sql.DataSource;
/**
 *  事務(wù)管理
 *  處理“只讀事務(wù)到讀庫,讀寫事務(wù)到寫庫”
 *  
 * @author sunchangtan 
 */
public class MasterSlaverDataSourceTransactionManager extends DataSourceTransactionManager {
    public MasterSlaverDataSourceTransactionManager(DataSource dataSource) {
        super(dataSource);
    }
    /**
     * 只讀事務(wù)到讀庫,讀寫事務(wù)到寫庫
     * @param transaction
     * @param definition
     */
    @Override
    protected void doBegin(Object transaction, TransactionDefinition definition) {
        //設(shè)置數(shù)據(jù)源
        boolean readOnly = definition.isReadOnly();
        if(readOnly) {
            DataSourceHolder.setCurrentDataSource(DataSourceType.READ);
        } else {
            DataSourceHolder.setCurrentDataSource(DataSourceType.WRITE);
        }
        super.doBegin(transaction, definition);
    }
    /**
     * 清理本地線程的數(shù)據(jù)源
     * @param transaction
     */
    @Override
    protected void doCleanupAfterCompletion(Object transaction) {
        super.doCleanupAfterCompletion(transaction);
        DataSourceHolder.clearDataSource();
    }
}

mybatis的讀寫分離的插件,需要配置到mybatis-config.xml

import lombok.extern.slf4j.Slf4j;
import org.apache.ibatis.executor.statement.RoutingStatementHandler;
import org.apache.ibatis.executor.statement.StatementHandler;
import org.apache.ibatis.logging.jdbc.ConnectionLogger;
import org.apache.ibatis.mapping.MappedStatement;
import org.apache.ibatis.plugin.*;
import org.apache.ibatis.reflection.DefaultReflectorFactory;
import org.apache.ibatis.reflection.MetaObject;
import org.apache.ibatis.reflection.factory.DefaultObjectFactory;
import org.apache.ibatis.reflection.wrapper.DefaultObjectWrapperFactory;
import org.springframework.jdbc.datasource.ConnectionProxy;
import java.lang.reflect.InvocationHandler;
import java.lang.reflect.Proxy;
import java.sql.Connection;
import java.util.Properties;
/**
 * 數(shù)據(jù)源讀寫分離路由
 */
@Slf4j
@Intercepts({@Signature(type = StatementHandler.class, method = "prepare", args = {Connection.class, Integer.class})})
public class MasterSlaveInterceptor implements Interceptor {
    public Object intercept(Invocation invocation) throws Throwable {
        Connection conn = (Connection) invocation.getArgs()[0];
        conn = unwrapConnection(conn);
        if (conn instanceof ConnectionProxy) {
            //強制走寫庫
            if (ConnectionHolder.FORCE_WRITE.get() != null && ConnectionHolder.FORCE_WRITE.get()) {
                if (log.isDebugEnabled()) {
                    log.debug("本事務(wù)強制走寫庫");
                }
                routeConnection(DataSourceType.WRITE, conn);
                return invocation.proceed();
            }
            StatementHandler statementHandler = (StatementHandler) invocation.getTarget();
            MetaObject metaObject = MetaObject.forObject(statementHandler, new DefaultObjectFactory(), new DefaultObjectWrapperFactory(), new DefaultReflectorFactory());
            MappedStatement mappedStatement;
            if (statementHandler instanceof RoutingStatementHandler) {
                mappedStatement = (MappedStatement) metaObject.getValue("delegate.mappedStatement");
            } else {
                mappedStatement = (MappedStatement) metaObject.getValue("mappedStatement");
            }
            if(mappedStatement.getId().endsWith(".insertSelective!selectKey")) {
                System.out.println("111");
            }
            DataSourceType key = DataSourceHolder.getCurrentDataSource();
            if (key == null) {
                key = DataSourceType.WRITE;
                String sel = statementHandler.getBoundSql().getSql().trim().substring(0, 3);
                if (sel.equalsIgnoreCase("sel")
                        && !mappedStatement.getId().endsWith(".insert!selectKey")
                        && !mappedStatement.getId().endsWith(".insertSelective!selectKey")) {
                    key = DataSourceType.READ;
                }
            }
            if(key == DataSourceType.WRITE) {
                if (log.isDebugEnabled()) {
                    log.debug("當(dāng)前數(shù)據(jù)庫為寫庫");
                }
            } else if(key == DataSourceType.READ) {
                if (log.isDebugEnabled()) {
                    log.debug("當(dāng)前數(shù)據(jù)庫為讀庫");
                }
            }
            routeConnection(key, conn);
        }
        return invocation.proceed();
    }
    private void routeConnection(DataSourceType key, Connection conn) {
        ConnectionHolder.CURRENT_CONNECTION.set(key);
        // 同一個線程下保證最多只有一個寫數(shù)據(jù)鏈接和讀數(shù)據(jù)鏈接
        if (!ConnectionHolder.CONNECTION_CONTEXT.get().containsKey(key)) {
            ConnectionProxy conToUse = (ConnectionProxy) conn;
            conn = conToUse.getTargetConnection();
            ConnectionHolder.CONNECTION_CONTEXT.get().put(key, conn);
        }
    }
    public Object plugin(Object target) {
        if (target instanceof StatementHandler) {
            return Plugin.wrap(target, this);
        } else {
            return target;
        }
    }
    public void setProperties(Properties properties) {
        // NOOP
    }
    /**
     * MyBatis wraps the JDBC Connection with a logging proxy but Spring registers the original connection so it should
     * be unwrapped before calling {@code DataSourceUtils.isConnectionTransactional(Connection, DataSource)}
     *
     * @param connection May be a {@code ConnectionLogger} proxy
     * @return the original JDBC {@code Connection}
     */
    private Connection unwrapConnection(Connection connection) {
        if (Proxy.isProxyClass(connection.getClass())) {
            InvocationHandler handler = Proxy.getInvocationHandler(connection);
            if (handler instanceof ConnectionLogger) {
                return ((ConnectionLogger) handler).getConnection();
            }
        }
        return connection;
    }
}

定義springboot的主從數(shù)據(jù)庫配置

import com.sample.dao.dynamic.DataSourceProxy;
import com.sample.dao.dynamic.DataSourceRouter;
import com.sample.dao.dynamic.RoundRobinMasterSlaverDataSourceRouter;
import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import javax.annotation.Resource;
import javax.sql.DataSource;
import java.util.Collections;
/**
 * 主從數(shù)據(jù)源的配置
 * @author sunchangtan
 */
@EnableConfigurationProperties({DataSourceMasterConfig.class, DataSourceSlaveConfig.class})
@Configuration
public class MasterSlaverDataSourceConfig {
    @Resource
    private DataSourceSlaveConfig dataSourceSlaveConfig;
    @Resource
    private DataSourceMasterConfig dataSourceMasterConfig;
    @Bean
    @ConditionalOnMissingBean
    public DataSourceRouter readRoutingDataSource() {
        RoundRobinMasterSlaverDataSourceRouter proxy = new RoundRobinMasterSlaverDataSourceRouter();
        proxy.setReadDataSoures(Collections.singletonList(dataSourceSlaveConfig.createDataSource()));
        proxy.setWriteDataSource(dataSourceMasterConfig.createDataSource());
        return proxy;
    }
    @Bean
    public DataSource dataSource(DataSourceRouter dataSourceRouter) {
        return new DataSourceProxy(dataSourceRouter);
    }
}

springboot中配置數(shù)據(jù)庫事務(wù)

import com.sample.dao.dynamic.MasterSlaverDataSourceTransactionManager;
import org.aspectj.lang.annotation.Aspect;
import org.springframework.aop.Advisor;
import org.springframework.aop.aspectj.AspectJExpressionPointcut;
import org.springframework.aop.support.DefaultPointcutAdvisor;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.transaction.PlatformTransactionManager;
import org.springframework.transaction.TransactionDefinition;
import org.springframework.transaction.interceptor.*;
import javax.annotation.Resource;
import javax.sql.DataSource;
import java.util.Collections;
import java.util.HashMap;
import java.util.Map;
/**
 * 事務(wù)管理的配置
 *
 * @author sunchangtan
 * @date 2018/8/30 11:22
 */
@Aspect
@Configuration
public class TransactionManagerConfigurer {
    private static final int TX_METHOD_TIMEOUT = 50000;
    private static final String AOP_POINTCUT_EXPRESSION = "execution(* com.sample.***.service..*.*(..))";
    @Resource
    private DataSource dataSource;
    @Bean
    public PlatformTransactionManager transactionManager() {
        return new MasterSlaverDataSourceTransactionManager(dataSource);
    }
    /**
     * 事務(wù)的實現(xiàn)Advice
     *
     * @return
     */
    @Bean
    public TransactionInterceptor txAdvice(PlatformTransactionManager transactionManager) {
        NameMatchTransactionAttributeSource source = new NameMatchTransactionAttributeSource();
        RuleBasedTransactionAttribute readOnlyTx = new RuleBasedTransactionAttribute();
        readOnlyTx.setReadOnly(true);
        //使用PROPAGATION_SUPPORTS:支持當(dāng)前事務(wù),如果當(dāng)前沒有事務(wù),就以非事務(wù)方式執(zhí)行。 如果查詢中出現(xiàn)異常,那么當(dāng)前事務(wù)也可以回滾
        readOnlyTx.setPropagationBehavior(TransactionDefinition.PROPAGATION_SUPPORTS);
        RuleBasedTransactionAttribute requiredTx = new RuleBasedTransactionAttribute();
        requiredTx.setRollbackRules(Collections.singletonList(new RollbackRuleAttribute(Exception.class)));
        //使用PROPAGATION_REQUIRED:如果當(dāng)前沒有事務(wù),就新建一個事務(wù),如果已經(jīng)存在一個事務(wù)中,加入到這個事務(wù)中。 如果需要數(shù)據(jù)庫增刪改,必須要使用事務(wù)
        requiredTx.setPropagationBehavior(TransactionDefinition.PROPAGATION_REQUIRED);
        requiredTx.setTimeout(TX_METHOD_TIMEOUT);
        Map<String, TransactionAttribute> txMap = new HashMap<>();
        txMap.put("add*", requiredTx);
        txMap.put("save*", requiredTx);
        txMap.put("insert*", requiredTx);
        txMap.put("update*", requiredTx);
        txMap.put("delete*", requiredTx);
        txMap.put("remove*", requiredTx);
        txMap.put("upload*", requiredTx);
        txMap.put("generate*", requiredTx);
        txMap.put("import*", requiredTx);
        txMap.put("bind*", requiredTx);
        txMap.put("unbind*", requiredTx);
        txMap.put("cancel*", requiredTx);
        txMap.put("send*", requiredTx);
        txMap.put("create*", requiredTx);
        txMap.put("compute*", requiredTx);
        txMap.put("recompute*", requiredTx);
        txMap.put("execute*", requiredTx);
        //txMap.put("submit*", requiredTx);
        txMap.put("get*", readOnlyTx);
        txMap.put("query*", readOnlyTx);
        txMap.put("list*", readOnlyTx);
        txMap.put("has*", readOnlyTx);
        txMap.put("exist*", readOnlyTx);
        txMap.put("download*", readOnlyTx);
        txMap.put("export*", readOnlyTx);
        txMap.put("search*", readOnlyTx);
        txMap.put("check*", readOnlyTx);
        txMap.put("load*", readOnlyTx);
        txMap.put("find*", readOnlyTx);
        source.setNameMap(txMap);
        return new TransactionInterceptor(transactionManager, source);
    }
    /**
     * 切面的定義,pointcut及advice
     *
     * @param txAdvice
     * @return
     */
    @Bean
    public Advisor txAdviceAdvisor(@Qualifier("txAdvice") TransactionInterceptor txAdvice) {
        AspectJExpressionPointcut pointcut = new AspectJExpressionPointcut();
        pointcut.setExpression(AOP_POINTCUT_EXPRESSION);
        return new DefaultPointcutAdvisor(pointcut, txAdvice);
    }
}

到此這篇關(guān)于Java基于SpringBoot和tk.mybatis實現(xiàn)事務(wù)讀寫分離代碼實例的文章就介紹到這了,更多相關(guān)SpringBoot事務(wù)讀寫分離實例內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • SpringBoot 整合 Shiro 密碼登錄的實現(xiàn)代碼

    SpringBoot 整合 Shiro 密碼登錄的實現(xiàn)代碼

    這篇文章主要介紹了SpringBoot 整合 Shiro 密碼登錄的實現(xiàn),本文通過實例代碼給大家介紹的非常詳細(xì),對大家的學(xué)習(xí)或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2021-02-02
  • IDEA中Maven Dependencies出現(xiàn)紅色波浪線的原因及解決方法

    IDEA中Maven Dependencies出現(xiàn)紅色波浪線的原因及解決方法

    在使用 IntelliJ IDEA 開發(fā) Java 項目時,尤其是基于 Maven 的項目,您可能會遇到 Maven Dependencies 中出現(xiàn)紅色波浪線的問題,這通常意味著項目的依賴無法被正確解析或下載,本文將詳細(xì)介紹該問題的原因及解決方法,并附有圖文說明,需要的朋友可以參考下
    2025-06-06
  • IDEA新建javaWeb以及Servlet簡單實現(xiàn)小結(jié)

    IDEA新建javaWeb以及Servlet簡單實現(xiàn)小結(jié)

    這篇文章主要介紹了IDEA新建javaWeb以及Servlet簡單實現(xiàn)小結(jié),小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2018-11-11
  • SpringBoot配置 Druid 三種方式(包括純配置文件配置)

    SpringBoot配置 Druid 三種方式(包括純配置文件配置)

    本文給大家分享在項目中用純 YML(application.yml 或者 application.properties)文件、Java 代碼配置 Bean 和注解三種方式配置 Alibaba Druid 用于監(jiān)控或者查看 SQL 狀況的相關(guān)知識,感興趣的朋友一起看看吧
    2021-10-10
  • java.io.UnsupportedEncodingException異常的正確解決方法(親測有效!)

    java.io.UnsupportedEncodingException異常的正確解決方法(親測有效!)

    這篇文章主要給大家介紹了關(guān)于java.io.UnsupportedEncodingException異常的正確解決方法,文中介紹的辦法親測有效,java.io.UnsupportedEncodingException是Java編程語言中的一個異常類,表示指定的字符集不被支持,需要的朋友可以參考下
    2024-02-02
  • 詳解java中finalize的實現(xiàn)與相應(yīng)的執(zhí)行過程

    詳解java中finalize的實現(xiàn)與相應(yīng)的執(zhí)行過程

    在常規(guī)的java書籍中,即會描述 object的finalize方法是用于一些特殊的對象在回收之前再做一些掃尾的工作,但是并沒有說明此是如何實現(xiàn)的.本篇從java的角度(不涉及jvm以及c++),有需要的朋友們可以參考借鑒。
    2016-09-09
  • 關(guān)于Java中增強for循環(huán)使用的注意事項

    關(guān)于Java中增強for循環(huán)使用的注意事項

    for循環(huán)語句是java循環(huán)語句中最常用的循環(huán)語句,一般用在循環(huán)次數(shù)已知的情況下使用,這篇文章主要給大家介紹了關(guān)于Java中增強for循環(huán)使用的注意事項,需要的朋友可以參考下
    2021-06-06
  • java web請求和響應(yīng)中出現(xiàn)中文亂碼問題的解析

    java web請求和響應(yīng)中出現(xiàn)中文亂碼問題的解析

    這篇文章主要為大家解析了java web請求和響應(yīng)中出現(xiàn)中文亂碼問題,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2016-10-10
  • 利用Java制作字符動畫實例代碼

    利用Java制作字符動畫實例代碼

    這篇文章主要給大家介紹了關(guān)于如何利用Java制作字符動畫的相關(guān)資料,文中通過示例代碼介紹的非常詳細(xì),對大家學(xué)習(xí)或者使用Java具有一定的參考學(xué)習(xí)價值,需要的朋友們下面來一起學(xué)習(xí)學(xué)習(xí)吧
    2019-05-05
  • Spring Boot 中的 @Field 注解的原理解析

    Spring Boot 中的 @Field 注解的原理解析

    本文詳細(xì)介紹了 Spring Boot 中的 @Field 注解的原理和使用方法,通過使用 @Field 注解,我們可以將 HTTP 請求中的參數(shù)值自動綁定到 Java 對象的屬性上,簡化了開發(fā)過程,提高了開發(fā)效率,感興趣的朋友跟隨小編一起看看吧
    2023-07-07

最新評論

荣昌县| 桐城市| 广州市| 兰州市| 广德县| 云安县| 运城市| 乐东| 平阳县| 手游| 和平县| 胶州市| 商河县| 奎屯市| 磐安县| 永州市| 元阳县| 平昌县| 双桥区| 阜康市| 荣昌县| 廊坊市| 雅江县| 通道| 仙游县| 尉犁县| 邛崃市| 休宁县| 晋中市| 蒲城县| 浠水县| 龙江县| 读书| 专栏| 托克逊县| 桐城市| 黔东| 雅安市| 阿城市| 繁峙县| 宣城市|