【问题标题】:Making massive amounts of individual row updates faster or more efficient使大量的单个行更新更快或更有效
【发布时间】:2014-11-06 15:33:44
【问题描述】:

我正在编写一个 java 应用程序,它将一个数据库的信息 (db2) 复制到另一个数据库 (sql server)。操作顺序很简单:

  1. 检查在特定时间范围内是否有任何更新
  2. 从第一个数据库中获取指定时间范围内的所有内容
  3. 将数据库信息映射到 POJO
  4. 将 POJO 的子集划分为线程(在属性文件中预定义 #)
  5. 线程在每个 POJO 中单独循环
  6. 更新第二个数据库

我的一切工作正常,但在一天中的某些时候,需要进行的更新量会大幅增加(可能达到数十万)。

您可以在下面看到我的代码的通用版本。它遵循应用程序的基本算法。对象是通用的,实际应用程序有 5 种不同类型的指定对象,每种对象都有自己的更新线程类。但下面的通用函数正是它们的样子。而在updateDatabase() 方法中,它们都被添加到threads 并同时运行。

private void updateDatabase()
{
    List<Thread> threads = new ArrayList<>();
    addObjectThreads( threads );     
    startThreads( threads );
    joinAllThreads( threads );
}

private void addObjectThreads( List<Thread> threads )
{
    List<Object> objects = getTransformService().getObjects();
    logger.info( "Found " + objects.size() + " Objects" );
    createThreads( threads, objects, ObjectUpdaterThread.class );
}

private void createThreads( List<Thread> threads, List<?> objects, Class threadClass )
{
    final int BASE_OBJECT_LOAD = 1;
    int objectLoad = objects.size() / Database.getMaxThreads() > 0 ? objects.size() / Database.getMaxThreads() + BASE_OBJECT_LOAD : BASE_OBJECT_LOAD;

    for (int i = 0; i < (objects.size() / objectLoad); ++i)
    {
        int startIndex = i * objectLoad;
        int endIndex = (i + 1) * objectLoad;
        try
        {
            List<?> objectSubList = objects.subList( startIndex, endIndex > objects.size() ? objects.size() : endIndex );
            threads.add( new Thread( (Thread) threadClass.getConstructor( List.class ).newInstance( objectSubList ) ) );
        }
        catch (Exception exception)
        {
            logger.error( exception.getMessage() );
        }
    }
}


public class ObjectUpdaterThread extends BaseUpdaterThread
{
    private List<Object> objects;
    final private Logger logger = Logger.getLogger( ObjectUpdaterThread.class );

    public ObjectUpdaterThread( List<Object> objects)
    {
        this.objects = objects;
    }

    public void run()
    {
        for (Object object : objects)
        {
            logger.info( "Now Updating Object: " + object.getId() );
            getTransformService().updateObject( object );
        }
    }
}

所有这些都转到类似于以下代码的 spring 服务。同样是通用的,但每种类型的对象都具有完全相同的逻辑类型。上面代码中的 getObjects() 只是传递给 DAO 的一行代码,因此无需真正发布。

@Service
@Scope(value = "prototype")
public class TransformServiceImpl implements TransformService
{
    final private Logger logger = Logger.getLogger( TransformServiceImpl.class );

    @Autowired
    private TransformDao transformDao;

    @Override
    public void updateObject( Object object )
    {
        String sql;
        if ( object.exists() )
        {
            sql = Object.Mapper.UPDATE;
        }
        else
        {
            sql = Object.Mapper.INSERT;
        }

        boolean isCompleted = false;
        while ( !isCompleted )
        {
            try
            {
                transformDao.updateObject( object, sql );
                isCompleted = true;
            }
            catch (Exception exception)
            {
                logger.error( exception.getMessage() );
                threadSleep();
                logger.info( "Now retrying update for Object: " + object.getId() );
            }
        }
        logger.info( "Updated Object: " + object.getId() );
    }
}

最后,这些都进入了如下所示的 DAO:

@Repository
@Scope(value = "prototype")
public class TransformDaoImpl implements TransformDao
{
    //@Resource is like @Autowired but with the added option of being able to specify the name
    //Good for autowiring two different instances of the same class [NamedParameterJdbcTemplate]
    //Another alternative = @Autowired @Qualifier(BEAN_NAME)
    @Resource(name = "db2")
    private NamedParameterJdbcTemplate db2;

    @Resource(name = "sqlServer")
    private NamedParameterJdbcTemplate sqlServer;

    final private Logger logger = Logger.getLogger( TransformerImpl.class );

    @Override
    public void updateObject( Objet object, String sql )
    {
        MapSqlParameterSource source = new MapSqlParameterSource();
        source.addValue( "column1_value", object.getColumn1Value() );
        //put all source values from the POJO in just like above

        sqlServer.update( sql, source );
    }
}

我的插入语句如下所示:

"INSERT INTO dbo.OBJECT_TABLE " +
"(COLUMN1, COLUMN2...) " +
"VALUES(:column1_value, :column2_value... "

我的更新语句如下所示:

"UPDATE dbo.OBJECT_TABLE SET " +
"COLUMN1 = :column1_value, COLUMN2 = :column2_value, " +
"WHERE PRIMARY_KEY_COLUMN = :primary_key_value"

它有很多我知道的代码和东西,但我只是想布置我所拥有的一切,希望我能得到帮助,使其更快或更高效。更新这么多行需要几个小时,如果只需要几个/几个小时而不是几个小时,那就太好了。谢谢你的帮助。我欢迎所有关于 Spring、线程和数据库的学习经验。

【问题讨论】:

  • 如果你被单独的行更新卡住了,它总是效率低下。您确定不能改用基于集合的逻辑吗?
  • 最有效的方法是转储表(并将它们导入目标)。如果您需要对已经存在的行的各个列进行更新,那您就很不走运了;如果您要更新“只是因为”,请尝试删除并重新插入。

标签: java sql-server multithreading spring transactions


【解决方案1】:

如果您要向服务器发送大量 SQL,您应该考虑使用 Statement.addBatchStatement.executeBatch 方法对它进行批处理。这些批次的大小是有限的(我总是将我的 SQL 限制为 64K),但它们大大降低了到数据库的往返次数。

当我迭代和创建 SQL 时,我会记录我已经批处理了多少,当 SQL 越过 64K 边界时,我会触发 executeBatch 并开始一个新的。

您可能想尝试使用 64K 数字,这可能是我当时使用的 Oracle 限制。

我无法与 Spring 交谈,但批处理是 JDBC Statement 的一部分。我确信这很简单。

【讨论】:

  • 你让我走上了正确的道路,谢谢。我使用了 NamedParameterJdbcTemplate 的 batchUpdate 方法,虽然仍然没有我想要的那么快,但它肯定快了一点。所以谢谢!
【解决方案2】:
  1. 检查在特定时间范围内是否有任何更新
  2. 从第一个数据库中获取指定时间范围内的所有内容

源表中的 LAST_UPDATED_DATE 列(或您正在使用的任何内容)上是否有索引?与其把负担放在您的应用程序上,如果它在您的控制范围内,为什么不在源数据库中编写一些触发器,在“更新日志”表中创建条目呢?这样,您的应用所需要做的就是使用并执行这些条目。

您如何管理您的交易?如果您为每个操作创建一个新事务,它将会非常缓慢。

关于线程代码,您是否考虑过使用更标准的代码而不是自己编写代码?你所拥有的是一个非常典型的producer/consumer,Java 对这种类型的东西有很好的支持,ThreadPoolExecutornumerous queue implementations 在执行不同任务的线程之间移动数据。

使用现成的东西的好处是 1) 它经过了很好的测试 2) 有许多调整选项和大小调整策略可供您调整以提高性能。

另外,您是否考虑过将每种类型的处理逻辑封装到单独的策略类中,而不是为需要处理的每种类型的对象使用 5 种不同的线程类型?这样,您可以使用单个工作线程池(这将更容易调整大小和调整)。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2012-08-16
    • 1970-01-01
    • 2016-10-01
    • 2013-08-22
    • 2016-08-19
    • 2015-11-30
    • 1970-01-01
    • 2014-10-12
    相关资源
    最近更新 更多