mongodb java 导入_亿级别记录的mongodb批量导入Es的java代码完整实现
importjava.io.IOException;importjava.net.UnknownHostException;importjava.util.HashMap;importjava.util.List;importjava.util.Map;importorg.apache.commons.codec.binary.Base64;importorg.apache.http.HttpHost;importorg.bson.types.ObjectId;importorg.elasticsearch.action.bulk.BulkItemResponse;importorg.elasticsearch.action.bulk.BulkRequest;importorg.elasticsearch.action.bulk.BulkResponse;importorg.elasticsearch.action.index.IndexRequest;importorg.elasticsearch.client.RestClient;importorg.elasticsearch.client.RestHighLevelClient;importorg.elasticsearch.common.xcontent.XContentType;importcom.mongodb.BasicDBObject;importcom.mongodb.DB;importcom.mongodb.DBCollection;importcom.mongodb.DBCursor;importcom.mongodb.DBObject;importcom.mongodb.MongoClient;importcom.mongodb.MongoException;public classTest {public static void main(String[] args) throwsIOException {int pageSize=10000;try{
MongoClient mongo= new MongoClient("localhost", 27017);/**** Get database ****/
//if database doesn't exists, MongoDB will create it for you
DB db = mongo.getDB("www");/**** Get collection / table from 'testdb' ****/
//if collection doesn't exists, MongoDB will create it for you
DBCollection table = db.getCollection("person");
RestHighLevelClient client= newRestHighLevelClient(
RestClient.builder(new HttpHost("localhost", 9200, "http")));
DBCursor dbObjects;
Long cnt=table.count();
System.out.println(table.getStats().toString());
Long page=getPageSize(cnt,pageSize);
ObjectId lastIdObject=null;
Long start=System.currentTimeMillis();long ss=start;for(Long i=0L;i
start=System.currentTimeMillis();
dbObjects=getCursorForCollection(table, lastIdObject, pageSize);
System.out.println("第"+(i+1)+"次查询,耗时:"+(System.currentTimeMillis()-start)+" 毫秒");
List objs=dbObjects.toArray();
start=System.currentTimeMillis();
batchInsertToEsSync(client,objs,"person","doc");
lastIdObject=(ObjectId) objs.get(objs.size()-1).get("_id");
System.out.println("第"+(i+1)+"次插入,耗时:"+(System.currentTimeMillis()-start)+" 毫秒");
}
System.out.println("耗时:"+(System.currentTimeMillis()-ss)/1000+"秒");
}catch(UnknownHostException e) {
e.printStackTrace();
}catch(MongoException e) {
e.printStackTrace();
}
}public static void batchInsertToEsSync(RestHighLevelClient client,List objs,String tableName,String type) throwsIOException {
BulkRequest bulkRequest=newBulkRequest();for(DBObject obj:objs) {
IndexRequest req= newIndexRequest(tableName, type);
Map map=new HashMap<>();for(String key:obj.keySet()) {if("_id".equalsIgnoreCase(key)) {
map.put("id", obj.get(key));
}else{
String valStr="";
Object val=obj.get(key);if(val!=null) {
valStr=Base64.encodeBase64String(val.toString().getBytes());
}
map.put(key, valStr);
}
}
req.id(map.get("id").toString());
req.source(map, XContentType.JSON);
bulkRequest.add(req);
}
BulkResponse bulkResponse=client.bulk(bulkRequest);for(BulkItemResponse bulkItemResponse : bulkResponse) {if(bulkItemResponse.isFailed()) {
System.out.println(bulkItemResponse.getId()+","+bulkItemResponse.getFailureMessage());
}
}
}public static DBCursor getCursorForCollection(DBCollection collection,ObjectId lastIdObject,intpageSize) {
DBCursor dbObjects=null;if(lastIdObject==null) {
lastIdObject=(ObjectId) collection.findOne().get("_id");
}
BasicDBObject query=newBasicDBObject();
query.append("_id",new BasicDBObject("$gt",lastIdObject));
BasicDBObject sort=newBasicDBObject();
sort.append("_id",1);
dbObjects=collection.find(query).limit(pageSize).sort(sort);returndbObjects;
}public static Long getPageSize(Long cnt,intpageSize) {return cnt%pageSize==0?cnt/pageSize:cnt/pageSize+1;
}
更多推荐
所有评论(0)