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

利用Java多線程技術(shù)導(dǎo)入數(shù)據(jù)到Elasticsearch的方法步驟

 更新時(shí)間:2019年07月16日 09:47:06   作者:Wooola  
這篇文章主要介紹了利用Java多線程技術(shù)導(dǎo)入數(shù)據(jù)到Elasticsearch的方法步驟,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧

前言


近期接到一個(gè)任務(wù),需要改造現(xiàn)有從mysql往Elasticsearch導(dǎo)入數(shù)據(jù)MTE(mysqlToEs)小工具,由于之前采用單線程導(dǎo)入,千億數(shù)據(jù)需要兩周左右的時(shí)間才能導(dǎo)入完成,導(dǎo)入效率非常低。所以樓主花了3天的時(shí)間,利用java線程池框架Executors中的FixedThreadPool線程池重寫了MTE導(dǎo)入工具,單臺(tái)服務(wù)器導(dǎo)入效率提高十幾倍(合理調(diào)整線程數(shù)據(jù),效率更高)。

關(guān)鍵技術(shù)棧

  • Elasticsearch
  • jdbc
  • ExecutorService\Thread
  • sql

工具說(shuō)明

maven依賴

<dependency> 
 <groupId>mysql</groupId> 
 <artifactId>mysql-connector-java</artifactId> 
 <version>${mysql.version}</version> 
</dependency> 
<dependency> 
 <groupId>org.elasticsearch</groupId> 
 <artifactId>elasticsearch</artifactId> 
 <version>${elasticsearch.version}</version> 
</dependency> 
<dependency> 
 <groupId>org.elasticsearch.client</groupId> 
 <artifactId>transport</artifactId> 
 <version>${elasticsearch.version}</version> 
</dependency> 
<dependency> 
 <groupId>org.projectlombok</groupId> 
 <artifactId>lombok</artifactId> 
 <version>${lombok.version}</version> 
</dependency> 
<dependency> 
 <groupId>com.alibaba</groupId> 
 <artifactId>fastjson</artifactId> 
 <version>${fastjson.version}</version> 
</dependency> 

java線程池設(shè)置

默認(rèn)線程池大小為21個(gè),可調(diào)整。其中POR為處理流程已辦數(shù)據(jù)線程池,ROR為處理流程已閱數(shù)據(jù)線程池。

private static int THREADS = 21; 
public static ExecutorService POR = Executors.newFixedThreadPool(THREADS); 
public static ExecutorService ROR = Executors.newFixedThreadPool(THREADS); 

定義已辦生產(chǎn)者線程/已閱生產(chǎn)者線程:ZlPendProducer/ZlReadProducer

public class ZlPendProducer implements Runnable { 
 ... 
 @Override 
 public void run() { 
 System.out.println(threadName + "::啟動(dòng)..."); 
 for (int j = 0; j < Const.TBL.TBL_PEND_COUNT; j++) 
 try { 
 .... 
 int size = 1000; 
 for (int i = 0; i < count; i += size) { 
 if (i + size > count) { 
 //作用為size最后沒有100條數(shù)據(jù)則剩余幾條newList中就裝幾條 
 size = count - i; 
 } 
 String sql = "select * from " + tableName + " limit " + i + ", " + size; 
 System.out.println(tableName + "::sql::" + sql); 
 rs = statement.executeQuery(sql); 
 List<HistPendingEntity> lst = new ArrayList<>(); 
 while (rs.next()) { 
 HistPendingEntity p = PendUtils.getHistPendingEntity(rs); 
 lst.add(p); 
 } 
 MteExecutor.POR.submit(new ZlPendConsumer(lst)); 
 Thread.sleep(2000); 
 } 
 .... 
 } catch (Exception e) { 
 e.printStackTrace(); 
 } 
 } 
} 
public class ZlReadProducer implements Runnable { 
 ...已閱生產(chǎn)者處理邏輯同已辦生產(chǎn)者 
} 

定義已辦消費(fèi)者線程/已閱生產(chǎn)者線程:ZlPendConsumer/ZlReadConsumer

public class ZlPendConsumer implements Runnable { 
 private String threadName; 
 private List<HistPendingEntity> lst; 
 public ZlPendConsumer(List<HistPendingEntity> lst) { 
 this.lst = lst; 
 } 
 @Override 
 public void run() { 
 ... 
 lst.forEach(v -> { 
 try { 
 String json = new Gson().toJson(v); 
 EsClient.addDataInJSON(json, Const.ES.HistPendDB_Index, Const.ES.HistPendDB_type, v.getPendingId(), null); 
 Const.COUNTER.LD_P.incrementAndGet(); 
 } catch (Exception e) { 
 e.printStackTrace(); 
 System.out.println("err::PendingId::" + v.getPendingId()); 
 } 
 }); 
 ... 
 } 
} 
public class ZlReadConsumer implements Runnable { 
 //已閱消費(fèi)者處理邏輯同已辦消費(fèi)者 
} 

定義導(dǎo)入Elasticsearch數(shù)據(jù)監(jiān)控線程:Monitor

監(jiān)控線程-Monitor為了計(jì)算每分鐘導(dǎo)入Elasticsearch的數(shù)據(jù)總條數(shù),利用監(jiān)控線程,可以調(diào)整線程池的線程數(shù)的大小,以便利用多線程更快速的導(dǎo)入數(shù)據(jù)。

public void monitorToES() { 
 new Thread(() -> { 
 while (true) { 
 StringBuilder sb = new StringBuilder(); 
 sb.append("已辦表數(shù)::").append(Const.TBL.TBL_PEND_COUNT) 
 .append("::已辦總數(shù)::").append(Const.COUNTER.LD_P_TOTAL) 
 .append("::已辦入庫(kù)總數(shù)::").append(Const.COUNTER.LD_P); 
 sb.append("~~~~已閱表數(shù)::").append(Const.TBL.TBL_READ_COUNT); 
 sb.append("::已閱總數(shù)::").append(Const.COUNTER.LD_R_TOTAL) 
 .append("::已閱入庫(kù)總數(shù)::").append(Const.COUNTER.LD_R); 
 if (ldPrevPendCount == 0 && ldPrevReadCount == 0) { 
 ldPrevPendCount = Const.COUNTER.LD_P.get(); 
 ldPrevReadCount = Const.COUNTER.LD_R.get(); 
 start = System.currentTimeMillis(); 
 } else { 
 long end = System.currentTimeMillis(); 
 if ((end - start) / 1000 >= 60) { 
 start = end; 
 sb.append("\n#########################################\n"); 
 sb.append("已辦每分鐘TPS::" + (Const.COUNTER.LD_P.get() - ldPrevPendCount) + "條"); 
 sb.append("::已閱每分鐘TPS::" + (Const.COUNTER.LD_R.get() - ldPrevReadCount) + "條"); 
 ldPrevPendCount = Const.COUNTER.LD_P.get(); 
 ldPrevReadCount = Const.COUNTER.LD_R.get(); 
 } 
 } 
 System.out.println(sb.toString()); 
 try { 
 Thread.sleep(3000); 
 } catch (InterruptedException e) { 
 e.printStackTrace(); 
 } 
 } 
 }).start(); 
} 

初始化Elasticsearch:EsClient

String cName = meta.get("cName");//es集群名字 
String esNodes = meta.get("esNodes");//es集群ip節(jié)點(diǎn) 
Settings esSetting = Settings.builder() 
 .put("cluster.name", cName) 
 .put("client.transport.sniff", true)//增加嗅探機(jī)制,找到ES集群 
 .put("thread_pool.search.size", 5)//增加線程池個(gè)數(shù),暫時(shí)設(shè)為5 
 .build(); 
String[] nodes = esNodes.split(","); 
client = new PreBuiltTransportClient(esSetting); 
for (String node : nodes) { 
 if (node.length() > 0) { 
 String[] hostPort = node.split(":"); 
 client.addTransportAddress(new TransportAddress(InetAddress.getByName(hostPort[0]), Integer.parseInt(hostPort[1]))); 
 } 
} 

初始化數(shù)據(jù)庫(kù)連接

conn = DriverManager.getConnection(url, user, password); 

啟動(dòng)參數(shù)

nohup java -jar mte.jar ES-Cluster2019 node1:9300,node2:9300,node3:9300 root 123456! jdbc:mysql://ip:3306/mte 130 130 >> ./mte.log 2>&1 & 

參數(shù)說(shuō)明

ES-Cluster2019 為Elasticsearch集群名字

node1:9300,node2:9300,node3:9300為es的節(jié)點(diǎn)IP

130 130為已辦已閱分表的數(shù)據(jù)

程序入口:MteMain

// 監(jiān)控線程 
Monitor monitorService = new Monitor(); 
monitorService.monitorToES(); 
// 已辦生產(chǎn)者線程 
Thread pendProducerThread = new Thread(new ZlPendProducer(conn, "ZlPendProducer")); 
pendProducerThread.start(); 
// 已閱生產(chǎn)者線程 
Thread readProducerThread = new Thread(new ZlReadProducer(conn, "ZlReadProducer")); 
readProducerThread.start(); 

以上就是本文的全部?jī)?nèi)容,希望對(duì)大家的學(xué)習(xí)有所幫助,也希望大家多多支持腳本之家。

相關(guān)文章

  • mybatis-plus如何禁用一級(jí)緩存的方法

    mybatis-plus如何禁用一級(jí)緩存的方法

    這篇文章主要介紹了mybatis-plus如何禁用一級(jí)緩存的方法,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧
    2021-03-03
  • java操作mysql實(shí)現(xiàn)增刪改查的方法

    java操作mysql實(shí)現(xiàn)增刪改查的方法

    這篇文章主要介紹了java操作mysql實(shí)現(xiàn)增刪改查的方法,結(jié)合實(shí)例形式分析了java操作mysql數(shù)據(jù)庫(kù)進(jìn)行增刪改查的具體實(shí)現(xiàn)技巧與相關(guān)注意事項(xiàng),需要的朋友可以參考下
    2017-05-05
  • Spring?@Scheduled定時(shí)器注解使用方式

    Spring?@Scheduled定時(shí)器注解使用方式

    這篇文章主要介紹了Spring?@Scheduled定時(shí)器注解使用方式,具有很好的參考價(jià)值,希望對(duì)大家有所幫助,如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2024-08-08
  • Java CAS操作與Unsafe類詳解

    Java CAS操作與Unsafe類詳解

    這篇文章主要介紹了Java CAS操作與Unsafe類的相關(guān)資料,幫助大家更好的理解和學(xué)習(xí)使用Java,感興趣的朋友可以了解下
    2021-02-02
  • SpringMVC請(qǐng)求流程源碼解析

    SpringMVC請(qǐng)求流程源碼解析

    這篇文章主要介紹了SpringMVC請(qǐng)求流程源碼分析,包括springmvc使用,SpringMVC啟動(dòng)過(guò)程及SpringMVC請(qǐng)求過(guò)程,本文通過(guò)實(shí)例代碼給大家介紹的非常詳細(xì),需要的朋友可以參考下
    2022-07-07
  • Java使用screw來(lái)對(duì)比數(shù)據(jù)庫(kù)表和字段差異

    Java使用screw來(lái)對(duì)比數(shù)據(jù)庫(kù)表和字段差異

    這篇文章主要介紹了Java如何使用screw來(lái)對(duì)比數(shù)據(jù)庫(kù)表和字段差異,文中的示例代碼講解詳細(xì),感興趣的小伙伴可以跟隨小編一起學(xué)習(xí)一下
    2024-12-12
  • Java 獲取Html文本中的img標(biāo)簽下src中的內(nèi)容方法

    Java 獲取Html文本中的img標(biāo)簽下src中的內(nèi)容方法

    今天小編就為大家分享一篇Java 獲取Html文本中的img標(biāo)簽下src中的內(nèi)容方法,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。一起跟隨小編過(guò)來(lái)看看吧
    2018-06-06
  • windows 32位eclipse遠(yuǎn)程hadoop開發(fā)環(huán)境搭建

    windows 32位eclipse遠(yuǎn)程hadoop開發(fā)環(huán)境搭建

    這篇文章主要介紹了windows 32位eclipse遠(yuǎn)程hadoop開發(fā)環(huán)境搭建的相關(guān)資料,需要的朋友可以參考下
    2016-07-07
  • Java xml數(shù)據(jù)格式返回實(shí)現(xiàn)操作

    Java xml數(shù)據(jù)格式返回實(shí)現(xiàn)操作

    這篇文章主要介紹了Java xml數(shù)據(jù)格式返回實(shí)現(xiàn)操作,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。一起跟隨小編過(guò)來(lái)看看吧
    2020-08-08
  • java實(shí)現(xiàn)頁(yè)面置換算法

    java實(shí)現(xiàn)頁(yè)面置換算法

    這篇文章主要為大家詳細(xì)介紹了java實(shí)現(xiàn)頁(yè)面置換算法,文中示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2020-08-08

最新評(píng)論

宁都县| 太康县| 遵化市| 红原县| 漾濞| 方正县| 礼泉县| 千阳县| 马龙县| 巫山县| 淄博市| 原平市| 香河县| 横山县| 闻喜县| 泽库县| 武威市| 文水县| 合阳县| 娄底市| 长春市| 潼关县| 木兰县| 水城县| 右玉县| 牙克石市| 收藏| 平舆县| 连南| 阳西县| 康乐县| 安义县| 伽师县| 保德县| 安阳市| 师宗县| 汉阴县| 泽普县| 澄城县| 彰化县| 榆树市|