【问题标题】:two threads querying a table两个线程查询一个表
【发布时间】:2014-04-10 08:49:21
【问题描述】:

我有一个大约 1 m 记录的巨大表,我想对所有记录进行一些处理,所以 1 线程方式,比如说... 1000 条记录,处理它们,再获取 1000 条记录等... 但是如果我想使用多任务处理呢?那是 2 个线程,每个线程获取 1000 条记录并并行处理,我如何确保每个线程将获取不同的 1000 条记录? 注意:我正在使用休眠

好像是这样的

public void run() {


    partList=getKParts(10);
    operateOnList(partList);


}

【问题讨论】:

  • 您可以在主线程中放置一个函数,该函数可由 2 个工作线程访问,这将为它们提供应获取的数据范围。例如,1call 将返回 0..9,第二个 10..19,无论是谁调用。

标签: java multithreading hibernate


【解决方案1】:

好的,你可以同步代码。

public class MyClass {

    private final HibernateFetcher hibernateFetcher = new HibernateFetcher();

    private class Worker implements Runnable {    
       public run() {    
         List partList = hibernateFetcher.fetchRecords();
         operateOnList(partList);    
       }
    }

    public void myBatchProcessor() {

      while(!hibernateFetcher.isFinished()) {
      // create *n* workers and go!

      }    
   }       
}

class HibernateFetcher {        

  private int count = 0; 
  private final Object lock = new Object();
  private volatile boolean isFinished = false;  
  public List fetchRecords() {

      Criteria criteria = ...;

      synchronized(lock) {
         criteria.setFirstResult(count) // offset
                 .setMaxResults(1000);
         count=count+1000;
      }
      List result = criteria.list();
      isFinished = result.length > 0 ? false: true;
      return result;
  }

  public synchronized boolean isFinished(){
    return isFinished;
  }

}

【讨论】:

  • 请检查语法,因为我没有测试过这段代码。我没有用于测试代码的休眠项目。
  • 嘿 Varun,你在哪里处理结果 (operateOnList) :)
  • 我之前没有使用过同步,我应该把它设为静态吗?
  • @Mak : 根据要求.. :)
  • 计数应该是不稳定的?和或静态要为所有线程共享对吗?
【解决方案2】:

如果我理解正确,您不希望预先获取 1m 条记录,而是希望以 1000 条为一批,然后在 2 个线程中处理它们,但使其并行。 首先,您必须使用 RowCount 或其他东西在数据库查询中实现分页类型功能。从 Java 中,您可以将 fromRowCount 传递给 toRowCount 并以 1000 个批次获取记录并在线程中并行处理它们。我在这里添加示例代码,但您必须进一步实现不同变量的逻辑。

        int totalRecordCount = 100000;
        int batchSize =1000;
        ExecutorService executor = Executors.newFixedThreadPool(totalRecordCount/batchSize);
        for(int x=0; x < totalRecordCount;){
            int toRowCount = x+batchSize;
            partList=getKParts(10,x,toRowCount);
            x= toRowCount + 1;
            executor.submit(new Runnable<>() {
                @Override
                public void run() {
                    operateOnList(partList);
                }
            });
        }

希望这会有所帮助。如果需要进一步澄清,请告诉我

【讨论】:

    【解决方案3】:

    如果您在数据库中的记录确实具有intlong 类型的主键,则为每个线程添加限制以仅从范围中获取记录:

    Thread1: 0000 - 0999, 2000 - 2999, etc
    Thread2: 1000 - 1999, 3000 - 3999, etc
    

    这样,每个线程只需要一个offset、一个counter 和一个increment。例如,Thread1 的偏移量为 0,而Thread2 的偏移量为 1000。由于此示例中有两个线程,因此您有 2000 的增量。对于每一轮,每个循环的计数器(从 0 开始)递增线程并计算下一个范围为:

    form = offset + (count * 2000) 到 = 从 + 999

    【讨论】:

    • 它的数据库排序
    • 没关系。只要可以订购主键,您就可以使用此策略。来自@Varun Achar 的那个应该也可以。
    【解决方案4】:
    import com.se.sas.persistance.utils.HibernateUtils;
    
    public class FinderWorker implements Runnable {
    
    
        @Override
        public void run() {
            operateOnList(getNParts(IndexLocker.getAllowedListSize()));
    
        }
    
        public List<Parts> getNParts(int listSize) {
    
            try {
    
                criteria = .....
                // *********** SYNCHRONIZATION OCCURS HERE ********************//
                criteria.setFirstResult(IndexLocker.getAvailableIndex());
                criteria.setMaxResults(listSize);
                partList = criteria.list();
    
            } catch (Exception e) {
                e.printStackTrace();
    
            } finally {
    
                session.close();
            }
            return partList;
        }
    
        public void operateOnList(List<Parts> partList) {
        ....
        }
    
    }
    

    储物柜类

    public class IndexLocker {
    
        private static AtomicInteger index = new AtomicInteger(0);
    
        private final static int batchSize = 1000;
    
        public IndexLocker() {
    
        }
    
        public static int getAllowedListSize() {
            return batchSize;
    
        }
    
        public static synchronized void incrmntIndex(int hop) {
            index.getAndAdd(hop);
        }
    
        public static synchronized int getAvailableIndex() {
    
            int result = index.get();
            index.getAndAdd(batchSize);
            return result;
        }
    
    }
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2019-08-10
      • 2013-10-13
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多