diff --git a/pom.xml b/pom.xml index b0792c2..83e59b9 100644 --- a/pom.xml +++ b/pom.xml @@ -1,418 +1,451 @@ - 4.0.0 - - 1.0.0 - java-study - 0.0.1-SNAPSHOT - jar - - java-study - http://maven.apache.org - - - UTF-8 - UTF-8 - 1.8 - 1.8 - 1.8 - - - - - - - - - junit - junit - 4.13.1 - - - - - - - - - org.slf4j - slf4j-api - 1.7.25 - - - - - ch.qos.logback - logback-classic - 1.2.3 - - - - ch.qos.logback - logback-core - 1.2.3 - - - - - - - - com.fasterxml.jackson.jaxrs - jackson-jaxrs-json-provider - 2.9.6 - - - - - com.google.code.gson - gson - 2.8.5 - - - - - com.alibaba - fastjson - 1.2.49 - - - - commons-codec - commons-codec - 1.11 - - - - - - commons-lang - commons-lang - 2.6 - - - - org.apache.commons - commons-lang3 - 3.7 - - - - - org.apache.commons - commons-compress - 1.18 - - - - - - - org.apache.poi - poi - 3.17 - - - - org.apache.poi - poi-ooxml - 3.17 - - - - - com.google.zxing - core - 3.0.0 - - - com.google.zxing - javase - 3.0.0 - - - - - - com.github.pagehelper - pagehelper - 4.1.0 - - - - - org.projectlombok - lombok - 1.16.18 - - - - - com.geccocrawler - gecco - 1.2.8 - - - - - - com.google.protobuf - protobuf-java - 3.5.1 - - - - - - org.quartz-scheduler - quartz - 2.3.0 - - - - org.quartz-scheduler - quartz-jobs - 2.3.0 - - - - org.springframework - spring-beans - 5.2.1.RELEASE - - - - - com.google.guava - guava - 27.1-jre - - - - - - - - - - - redis.clients - jedis - 2.9.0 - - - - - mysql - mysql-connector-java - 5.1.41 - - - - - org.xerial - sqlite-jdbc - 3.20.1 - - - - - com.microsoft.sqlserver - sqljdbc4 - 4.0 - - - - - - - - - - - io.netty - netty-all - 4.1.42.Final - - - - - org.apache.mina - mina-core - 2.0.16 - - - - - com.github.kevinsawicki - http-request - 6.0 - - - - - - - - - - - - - com.rabbitmq - amqp-client - 4.3.0 - - - - - org.apache.kafka - kafka_2.12 - 1.0.0 - - - org.apache.zookeeper - zookeeper - - - org.slf4j - slf4j-log4j12 - - - log4j - log4j - - - - - - - org.apache.kafka - kafka-clients - 1.0.0 - - - - org.apache.kafka - kafka-streams - 1.0.0 - - - - - - - - - - org.apache.zookeeper - zookeeper - 3.4.10 - - - org.slf4j - slf4j-log4j12 - - - log4j - log4j - - - - - - - com.101tec - zkclient - 0.10 - - - - - org.apache.hbase - hbase-client - 2.1.4 - - - org.apache.hbase - hbase-common - 2.1.4 - - - - org.apache.hadoop - hadoop-client - 3.2.0 - - - - org.apache.spark - spark-core_2.12 - 2.4.0 - - - - org.apache.spark - spark-sql_2.12 - 2.4.0 - - - - - org.apache.storm - storm-core - 1.2.2 - provided - - - - org.apache.storm - storm-kafka - 1.2.3 - provided - - - - - - - - - org.elasticsearch.client - elasticsearch-rest-high-level-client - 6.6.1 - - - - - io.searchbox - jest - 6.3.1 - - - - - com.alibaba - easyexcel - 2.2.7 - - - - - - - - - - spring-milestone - http://repo.spring.io/libs-release - - + xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" + xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> + 4.0.0 + + 1.0.0 + java-study + 0.0.1-SNAPSHOT + jar + + java-study + http://maven.apache.org + + + UTF-8 + UTF-8 + 1.8 + 1.8 + 1.8 + + + + + + org.apache.maven.plugins + maven-compiler-plugin + 3.6.1 + + 1.8 + 1.8 + + + + + + + + + + + junit + junit + 4.13.1 + + + + + + + + + org.slf4j + slf4j-api + 1.7.25 + + + + + ch.qos.logback + logback-classic + 1.2.3 + + + + ch.qos.logback + logback-core + 1.2.3 + + + + + + + + com.fasterxml.jackson.jaxrs + jackson-jaxrs-json-provider + 2.9.6 + + + + + com.google.code.gson + gson + 2.8.5 + + + + + com.alibaba + fastjson + 1.2.49 + + + + commons-codec + commons-codec + 1.11 + + + + + + commons-lang + commons-lang + 2.6 + + + + org.apache.commons + commons-lang3 + 3.7 + + + + + org.apache.commons + commons-compress + 1.18 + + + + + + org.apache.poi + poi + 3.17 + + + + org.apache.poi + poi-ooxml + 3.17 + + + + + com.google.zxing + core + 3.0.0 + + + com.google.zxing + javase + 3.0.0 + + + + + + com.github.pagehelper + pagehelper + 4.1.0 + + + + + org.projectlombok + lombok + 1.16.18 + + + + + com.geccocrawler + gecco + 1.2.8 + + + + + + com.google.protobuf + protobuf-java + 3.5.1 + + + + + + org.quartz-scheduler + quartz + 2.3.0 + + + + org.quartz-scheduler + quartz-jobs + 2.3.0 + + + + org.springframework + spring-beans + 5.2.1.RELEASE + + + + + com.google.guava + guava + 27.1-jre + + + + + + + + + + redis.clients + jedis + 2.9.0 + + + + + mysql + mysql-connector-java + 5.1.41 + + + + + org.xerial + sqlite-jdbc + 3.20.1 + + + + + com.microsoft.sqlserver + sqljdbc4 + 4.0 + + + + + + + + + + io.netty + netty-all + 4.1.42.Final + + + + + org.apache.mina + mina-core + 2.0.16 + + + + + com.github.kevinsawicki + http-request + 6.0 + + + + + + + + + + + com.rabbitmq + amqp-client + 4.3.0 + + + + org.apache.rocketmq + rocketmq-client + 4.7.1 + + + + + org.apache.kafka + kafka_2.12 + 1.0.0 + + + org.apache.zookeeper + zookeeper + + + org.slf4j + slf4j-log4j12 + + + log4j + log4j + + + + + + + org.apache.kafka + kafka-clients + 1.0.0 + + + + org.apache.kafka + kafka-streams + 1.0.0 + + + + + + + + + + org.apache.zookeeper + zookeeper + 3.4.10 + + + org.slf4j + slf4j-log4j12 + + + log4j + log4j + + + + + + + com.101tec + zkclient + 0.10 + + + + + org.apache.hbase + hbase-client + 2.1.4 + + + org.apache.hbase + hbase-common + 2.1.4 + + + + org.apache.hadoop + hadoop-client + 3.2.0 + + + + org.apache.spark + spark-core_2.12 + 2.4.0 + + + + org.apache.spark + spark-sql_2.12 + 2.4.0 + + + + + org.apache.storm + storm-core + 1.2.2 + provided + + + + org.apache.storm + storm-kafka + 1.2.3 + provided + + + + + + + + + + + + + + + org.elasticsearch.client + elasticsearch-rest-high-level-client + 7.7.1 + + + org.elasticsearch.client + elasticsearch-rest-client + 7.7.1 + + + org.elasticsearch + elasticsearch + 7.7.1 + + + + + io.searchbox + jest + 6.3.1 + + + + + com.alibaba + easyexcel + 2.2.7 + + + org.projectlombok + lombok + 1.18.24 + + + + + + + spring-milestone + http://repo.spring.io/libs-release + + + diff --git a/src/main/java/com/pancm/elasticsearch/EsAggregationSearchTest.java b/src/main/java/com/pancm/elasticsearch/EsAggregationSearchTest.java index 22b3049..f6576fa 100644 --- a/src/main/java/com/pancm/elasticsearch/EsAggregationSearchTest.java +++ b/src/main/java/com/pancm/elasticsearch/EsAggregationSearchTest.java @@ -23,17 +23,17 @@ import org.elasticsearch.search.aggregations.AggregationBuilder; import org.elasticsearch.search.aggregations.AggregationBuilders; import org.elasticsearch.search.aggregations.Aggregations; +import org.elasticsearch.search.aggregations.PipelineAggregatorBuilders; import org.elasticsearch.search.aggregations.bucket.histogram.DateHistogramInterval; import org.elasticsearch.search.aggregations.bucket.terms.Terms; import org.elasticsearch.search.aggregations.bucket.terms.TermsAggregationBuilder; -import org.elasticsearch.search.aggregations.metrics.avg.Avg; -import org.elasticsearch.search.aggregations.metrics.cardinality.CardinalityAggregationBuilder; -import org.elasticsearch.search.aggregations.metrics.max.Max; -import org.elasticsearch.search.aggregations.metrics.min.Min; -import org.elasticsearch.search.aggregations.metrics.sum.Sum; -import org.elasticsearch.search.aggregations.metrics.tophits.TopHits; -import org.elasticsearch.search.aggregations.pipeline.PipelineAggregatorBuilders; -import org.elasticsearch.search.aggregations.pipeline.bucketselector.BucketSelectorPipelineAggregationBuilder; +import org.elasticsearch.search.aggregations.metrics.Avg; +import org.elasticsearch.search.aggregations.metrics.CardinalityAggregationBuilder; +import org.elasticsearch.search.aggregations.metrics.Max; +import org.elasticsearch.search.aggregations.metrics.Min; +import org.elasticsearch.search.aggregations.metrics.Sum; +import org.elasticsearch.search.aggregations.metrics.TopHits; +import org.elasticsearch.search.aggregations.pipeline.BucketSelectorPipelineAggregationBuilder; import org.elasticsearch.search.builder.SearchSourceBuilder; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -58,9 +58,10 @@ public class EsAggregationSearchTest { + private static String elasticIp = "192.168.77.130"; - private static String elasticIp = "192.169.0.23"; private static int elasticPort = 9200; + private static Logger logger = LoggerFactory.getLogger(EsHighLevelRestSearchTest.class); private static RestHighLevelClient client = null; @@ -84,7 +85,7 @@ public static void main(String[] args) { topSearch(); } catch (Exception e) { e.printStackTrace(); - }finally { + } finally { close(); } @@ -108,8 +109,8 @@ private static void close() { client.close(); } catch (IOException e) { e.printStackTrace(); - }finally{ - client=null; + } finally { + client = null; } } } @@ -166,8 +167,8 @@ private static void createIndex() throws IOException { getRequest.humanReadable(true); boolean exists2 = client.indices().exists(getRequest, RequestOptions.DEFAULT); //如果存在就不创建了 - if(exists2) { - System.out.println(index+"索引库已经存在!"); + if (exists2) { + System.out.println(index + "索引库已经存在!"); return; } // 开始创建库 @@ -181,8 +182,8 @@ private static void createIndex() throws IOException { request.alias(new Alias("pancm_alias")); CreateIndexResponse createIndexResponse = client.indices().create(request, RequestOptions.DEFAULT); boolean falg = createIndexResponse.isAcknowledged(); - if(falg){ - System.out.println("创建索引库:"+index+"成功!" ); + if (falg) { + System.out.println("创建索引库:" + index + "成功!"); } } catch (IOException e) { e.printStackTrace(); @@ -191,39 +192,38 @@ private static void createIndex() throws IOException { } - /** * 批量操作示例 * * @throws InterruptedException */ - private static void bulk() throws IOException{ + private static void bulk() throws IOException { // 类型 String type = "_doc"; String index = "student"; BulkRequest request = new BulkRequest(); - int k =10; - List> mapList = new ArrayList<>(); + int k = 10; + List> mapList = new ArrayList<>(); LocalDateTime ldt = LocalDateTime.now(); - for (int i = 1; i <=k ; i++) { - Map map = new HashMap<>(); - map.put("uid",i); - map.put("age",i); - map.put("name","虚无境"+(i%3)); - map.put("class",i%10); - map.put("grade",400+i); - map.put("createtm",ldt.plusDays(i).format(DateTimeFormatter.ofPattern("yyyy-MM-dd"))); - map.put("updatetm",ldt.plusDays(i).format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss.SSS"))); - if(i==5){ - map.put("updatetm","2019-11-31 21:04:55.268"); + for (int i = 1; i <= k; i++) { + Map map = new HashMap<>(); + map.put("uid", i); + map.put("age", i); + map.put("name", "虚无境" + (i % 3)); + map.put("class", i % 10); + map.put("grade", 400 + i); + map.put("createtm", ldt.plusDays(i).format(DateTimeFormatter.ofPattern("yyyy-MM-dd"))); + map.put("updatetm", ldt.plusDays(i).format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss.SSS"))); + if (i == 5) { + map.put("updatetm", "2019-11-31 21:04:55.268"); } mapList.add(map); } - for (int i = 0; i map = mapList.get(i); + for (int i = 0; i < mapList.size(); i++) { + Map map = mapList.get(i); String id = map.get("uid").toString(); // 可以进行修改/删除/新增 操作 //docAsUpsert 为true表示存在更新,不存在插入,为false表示不存在就是不做更新 @@ -235,15 +235,15 @@ private static void bulk() throws IOException{ } /** + * @return void * @Author pancm * @Description 多个聚合条件测试 * SQL: select age, name, count(*) as count1 from student group by age, name; - * @Date 2019/7/3 + * @Date 2019/7/3 * @Param [] - * @return void **/ - private static void groupbySearch() throws IOException{ - String buk="group"; + private static void groupbySearch() throws IOException { + String buk = "group"; AggregationBuilder aggregation = AggregationBuilders.terms("age").field("age"); AggregationBuilder aggregation2 = AggregationBuilders.terms("name").field("name"); //根据创建时间按天分组 @@ -254,146 +254,146 @@ private static void groupbySearch() throws IOException{ aggregation2.subAggregation(aggregation3); aggregation.subAggregation(aggregation2); - agg(aggregation,buk); + agg(aggregation, buk); } /** + * @return void * @Author pancm * @Description 平均聚合查询测试用例 - * @Date 2019/4/1 + * @Date 2019/4/1 * @Param [] - * @return void **/ - private static void avgSearch() throws IOException { + private static void avgSearch() throws IOException { - String buk="t_grade_avg"; + String buk = "t_grade_avg"; //直接求平均数 AggregationBuilder aggregation = AggregationBuilders.avg(buk).field("grade"); logger.info("求班级的平均分数:"); - agg(aggregation,buk); + agg(aggregation, buk); } - private static void maxSearch() throws IOException{ - String buk="t_grade"; + private static void maxSearch() throws IOException { + String buk = "t_grade"; AggregationBuilder aggregation = AggregationBuilders.max(buk).field("grade"); logger.info("求班级的最高分数:"); - agg(aggregation,buk); + agg(aggregation, buk); } - private static void sumSearch() throws IOException{ - String buk="t_grade"; + private static void sumSearch() throws IOException { + String buk = "t_grade"; AggregationBuilder aggregation = AggregationBuilders.sum(buk).field("grade"); logger.info("求班级的总分数:"); - agg(aggregation,buk); + agg(aggregation, buk); } /** + * @return void * @Author pancm * @Description 平均聚合查询测试用例 - * @Date 2019/4/1 + * @Date 2019/4/1 * @Param [] - * @return void **/ - private static void avgGroupSearch() throws IOException { + private static void avgGroupSearch() throws IOException { - String agg="t_class_avg"; - String buk="t_grade"; + String agg = "t_class_avg"; + String buk = "t_grade"; //terms 就是分组统计 根据student的grade成绩进行分组并创建一个新的聚合 TermsAggregationBuilder aggregation = AggregationBuilders.terms(agg).field("class"); aggregation.subAggregation(AggregationBuilders.avg(buk).field("grade")); logger.info("根据班级求平均分数:"); - agg(aggregation,agg,buk); + agg(aggregation, agg, buk); } - private static void maxGroupSearch() throws IOException{ + private static void maxGroupSearch() throws IOException { - String agg="t_class_max"; - String buk="t_grade"; + String agg = "t_class_max"; + String buk = "t_grade"; //terms 就是分组统计 根据student的grade成绩进行分组并创建一个新的聚合 TermsAggregationBuilder aggregation = AggregationBuilders.terms(agg).field("class"); aggregation.subAggregation(AggregationBuilders.max(buk).field("grade")); logger.info("根据班级求最大分数:"); - agg(aggregation,agg,buk); + agg(aggregation, agg, buk); } - private static void sumGroupSearch() throws IOException{ - String agg="t_class_sum"; - String buk="t_grade"; + private static void sumGroupSearch() throws IOException { + String agg = "t_class_sum"; + String buk = "t_grade"; //terms 就是分组统计 根据student的grade成绩进行分组并创建一个新的聚合 TermsAggregationBuilder aggregation = AggregationBuilders.terms(agg).field("class"); aggregation.subAggregation(AggregationBuilders.sum(buk).field("grade")); logger.info("根据班级求总分:"); - agg(aggregation,agg,buk); + agg(aggregation, agg, buk); } - protected static void agg(AggregationBuilder aggregation, String buk) throws IOException{ + protected static void agg(AggregationBuilder aggregation, String buk) throws IOException { SearchResponse searchResponse = search(aggregation); - if(RestStatus.OK.equals(searchResponse.status())) { + if (RestStatus.OK.equals(searchResponse.status())) { // 获取聚合结果 Aggregations aggregations = searchResponse.getAggregations(); - if(buk.contains("avg")){ + if (buk.contains("avg")) { //取子聚合 Avg ba = aggregations.get(buk); - logger.info(buk+":" + ba.getValue()); + logger.info(buk + ":" + ba.getValue()); logger.info("------------------------------------"); - }else if(buk.contains("max")){ + } else if (buk.contains("max")) { //取子聚合 Max ba = aggregations.get(buk); - logger.info(buk+":" + ba.getValue()); + logger.info(buk + ":" + ba.getValue()); logger.info("------------------------------------"); - }else if(buk.contains("min")){ + } else if (buk.contains("min")) { //取子聚合 Min ba = aggregations.get(buk); - logger.info(buk+":" + ba.getValue()); + logger.info(buk + ":" + ba.getValue()); logger.info("------------------------------------"); - }else if(buk.contains("sum")){ + } else if (buk.contains("sum")) { //取子聚合 Sum ba = aggregations.get(buk); - logger.info(buk+":" + ba.getValue()); + logger.info(buk + ":" + ba.getValue()); logger.info("------------------------------------"); - }else if(buk.contains("top")){ + } else if (buk.contains("top")) { //取子聚合TopHits TopHits ba = aggregations.get(buk); - logger.info(buk+":" + ba.getHits().totalHits); + logger.info(buk + ":" + ba.getHits().getTotalHits()); logger.info("------------------------------------"); - }else if (buk.contains("group")){ - Map map = new HashMap<>(); - List> list = new ArrayList<>(); - agg(map,list,aggregations); - logger.info("聚合查询结果:"+list); + } else if (buk.contains("group")) { + Map map = new HashMap<>(); + List> list = new ArrayList<>(); + agg(map, list, aggregations); + logger.info("聚合查询结果:" + list); logger.info("------------------------------------"); } } } - private static void agg(Map map, List> list, Aggregations aggregations) { + private static void agg(Map map, List> list, Aggregations aggregations) { aggregations.forEach(aggregation -> { String name = aggregation.getName(); Terms genders = aggregations.get(name); for (Terms.Bucket entry : genders.getBuckets()) { String key = entry.getKey().toString(); long t = entry.getDocCount(); - map.put(name,key); - map.put(name+"_"+"count",t); + map.put(name, key); + map.put(name + "_" + "count", t); //判断里面是否还有嵌套的数据 List list2 = entry.getAggregations().asList(); if (list2.isEmpty()) { - Map map2 = new HashMap<>(); - BeanUtils.copyProperties(map,map2); + Map map2 = new HashMap<>(); + BeanUtils.copyProperties(map, map2); list.add(map2); - }else{ + } else { agg(map, list, entry.getAggregations()); } } @@ -407,14 +407,14 @@ private static void agg(List> list, Aggregations aggregation for (Terms.Bucket entry : genders.getBuckets()) { String key = entry.getKey().toString(); long t = entry.getDocCount(); - Map map =new HashMap<>(); - map.put(name,key); - map.put(name+"_"+"count",t); + Map map = new HashMap<>(); + map.put(name, key); + map.put(name + "_" + "count", t); //判断里面是否还有嵌套的数据 List list2 = entry.getAggregations().asList(); if (list2.isEmpty()) { list.add(map); - }else{ + } else { agg(list, entry.getAggregations()); } } @@ -423,7 +423,6 @@ private static void agg(List> list, Aggregations aggregation } - private static SearchResponse search(AggregationBuilder aggregation) throws IOException { SearchRequest searchRequest = new SearchRequest(); searchRequest.indices("student"); @@ -436,66 +435,66 @@ private static SearchResponse search(AggregationBuilder aggregation) throws IOEx //不需要版本号 searchSourceBuilder.version(false); searchSourceBuilder.aggregation(aggregation); - logger.info("查询的语句:"+searchSourceBuilder.toString()); + logger.info("查询的语句:" + searchSourceBuilder.toString()); searchRequest.source(searchSourceBuilder); // 同步查询 SearchResponse searchResponse = client.search(searchRequest, RequestOptions.DEFAULT); - return searchResponse; + return searchResponse; } /** + * @return void * @Author pancm * @Description 进行聚合 - * @Date 2019/4/2 + * @Date 2019/4/2 * @Param [] - * @return void **/ - protected static void agg(AggregationBuilder aggregation, String agg, String buk) throws IOException{ + protected static void agg(AggregationBuilder aggregation, String agg, String buk) throws IOException { // 同步查询 SearchResponse searchResponse = search(aggregation); //4、处理响应 //搜索结果状态信息 - if(RestStatus.OK.equals(searchResponse.status())) { + if (RestStatus.OK.equals(searchResponse.status())) { // 获取聚合结果 Aggregations aggregations = searchResponse.getAggregations(); //分组 Terms byAgeAggregation = aggregations.get(agg); - logger.info(agg+" 结果"); + logger.info(agg + " 结果"); logger.info("name: " + byAgeAggregation.getName()); logger.info("type: " + byAgeAggregation.getType()); logger.info("sumOfOtherDocCounts: " + byAgeAggregation.getSumOfOtherDocCounts()); logger.info("------------------------------------"); - for(Terms.Bucket buck : byAgeAggregation.getBuckets()) { + for (Terms.Bucket buck : byAgeAggregation.getBuckets()) { logger.info("key: " + buck.getKeyAsNumber()); logger.info("docCount: " + buck.getDocCount()); logger.info("docCountError: " + buck.getDocCountError()); - if(agg.contains("avg")){ + if (agg.contains("avg")) { //取子聚合 Avg ba = buck.getAggregations().get(buk); - logger.info(buk+":" + ba.getValue()); + logger.info(buk + ":" + ba.getValue()); logger.info("------------------------------------"); - }else if(agg.contains("max")){ + } else if (agg.contains("max")) { //取子聚合 Max ba = buck.getAggregations().get(buk); - logger.info(buk+":" + ba.getValue()); + logger.info(buk + ":" + ba.getValue()); logger.info("------------------------------------"); - }else if(agg.contains("min")){ + } else if (agg.contains("min")) { //取子聚合 Min ba = buck.getAggregations().get(buk); - logger.info(buk+":" + ba.getValue()); + logger.info(buk + ":" + ba.getValue()); logger.info("------------------------------------"); - }else if(agg.contains("sum")){ + } else if (agg.contains("sum")) { //取子聚合 Sum ba = buck.getAggregations().get(buk); - logger.info(buk+":" + ba.getValue()); + logger.info(buk + ":" + ba.getValue()); logger.info("------------------------------------"); } } @@ -503,14 +502,14 @@ protected static void agg(AggregationBuilder aggregation, String agg, String b } /** + * @return void * @Author pancm * @Description having - * @Date 2020/8/21 + * @Date 2020/8/21 * @Param [] - * @return void **/ - private static void havingSearch() throws IOException{ - String index=""; + private static void havingSearch() throws IOException { + String index = ""; SearchRequest searchRequest = new SearchRequest(index); searchRequest.indices(index); SearchSourceBuilder sourceBuilder = new SearchSourceBuilder(); @@ -520,15 +519,15 @@ private static void havingSearch() throws IOException{ String group_name = "nas_ip_address"; String query_name = "acct_start_time"; String query_type = "gte,lte"; - String query_name_value="2020-08-05 13:25:55,2020-08-20 13:26:55"; - String[] query_types= query_type.split(","); - String[] query_name_values= query_name_value.split(","); + String query_name_value = "2020-08-05 13:25:55,2020-08-20 13:26:55"; + String[] query_types = query_type.split(","); + String[] query_name_values = query_name_value.split(","); for (int i = 0; i < query_types.length; i++) { - if("gte".equals(query_types[i])){ + if ("gte".equals(query_types[i])) { boolQueryBuilder.must(QueryBuilders.rangeQuery(query_name).gte(query_name_values[i])); } - if("lte".equals(query_types[i])){ + if ("lte".equals(query_types[i])) { boolQueryBuilder.must(QueryBuilders.rangeQuery(query_name).lte(query_name_values[i])); } } @@ -562,26 +561,26 @@ private static void havingSearch() throws IOException{ long count = searchResponse.getHits().getHits().length; Aggregations aggregations = searchResponse.getAggregations(); // agg(aggregations); - Map map =new HashMap<>(); - List> list =new ArrayList<>(); - agg(list,aggregations); + Map map = new HashMap<>(); + List> list = new ArrayList<>(); + agg(list, aggregations); // System.out.println(map); System.out.println(list); } /** + * @return void * @Author pancm * @Description 去重 - * @Date 2020/8/26 + * @Date 2020/8/26 * @Param [] - * @return void **/ - private static void distinctSearch() throws IOException{ - String buk="group"; - String distinctName="name"; + private static void distinctSearch() throws IOException { + String buk = "group"; + String distinctName = "name"; AggregationBuilder aggregation = AggregationBuilders.terms("age").field("age"); - CardinalityAggregationBuilder cardinalityBuilder = AggregationBuilders.cardinality(distinctName).field(distinctName); + CardinalityAggregationBuilder cardinalityBuilder = AggregationBuilders.cardinality(distinctName).field(distinctName); //根据创建时间按天分组 // AggregationBuilder aggregation3 = AggregationBuilders.dateHistogram("createtm") // .field("createtm") @@ -590,12 +589,11 @@ private static void distinctSearch() throws IOException{ // // aggregation2.subAggregation(aggregation3); aggregation.subAggregation(cardinalityBuilder); - agg(aggregation,buk); + agg(aggregation, buk); } - - private static void topSearch() throws IOException{ + private static void topSearch() throws IOException { } diff --git a/src/main/java/com/pancm/elasticsearch/EsHighLevelRestSearchTest.java b/src/main/java/com/pancm/elasticsearch/EsHighLevelRestSearchTest.java index 0ef3d42..a59e725 100644 --- a/src/main/java/com/pancm/elasticsearch/EsHighLevelRestSearchTest.java +++ b/src/main/java/com/pancm/elasticsearch/EsHighLevelRestSearchTest.java @@ -1,6 +1,7 @@ package com.pancm.elasticsearch; import org.apache.http.HttpHost; +import org.apache.lucene.search.TotalHits; import org.elasticsearch.action.search.SearchRequest; import org.elasticsearch.action.search.SearchResponse; import org.elasticsearch.action.search.ShardSearchFailure; @@ -22,7 +23,7 @@ import org.elasticsearch.search.aggregations.bucket.terms.Terms; import org.elasticsearch.search.aggregations.bucket.terms.Terms.Bucket; import org.elasticsearch.search.aggregations.bucket.terms.TermsAggregationBuilder; -import org.elasticsearch.search.aggregations.metrics.avg.Avg; +import org.elasticsearch.search.aggregations.metrics.Avg; import org.elasticsearch.search.builder.SearchSourceBuilder; import org.elasticsearch.search.fetch.subphase.highlight.HighlightBuilder; import org.elasticsearch.search.suggest.Suggest; @@ -49,7 +50,7 @@ */ public class EsHighLevelRestSearchTest { - private static String elasticIp = "192.169.0.23"; + private static String elasticIp = "192.168.77.130"; private static int elasticPort = 9200; private static Logger logger = LoggerFactory.getLogger(EsHighLevelRestSearchTest.class); @@ -485,7 +486,7 @@ private static void search() throws IOException { SearchHits hits = searchResponse2.getHits(); //总条数和分值 - long totalHits = hits.getTotalHits(); + TotalHits totalHits = hits.getTotalHits(); float maxScore = hits.getMaxScore(); diff --git a/src/main/java/com/pancm/elasticsearch/EsHighLevelRestTest1.java b/src/main/java/com/pancm/elasticsearch/EsHighLevelRestTest1.java index 0238621..0c99203 100644 --- a/src/main/java/com/pancm/elasticsearch/EsHighLevelRestTest1.java +++ b/src/main/java/com/pancm/elasticsearch/EsHighLevelRestTest1.java @@ -60,7 +60,7 @@ */ public class EsHighLevelRestTest1 { - private static String elasticIp = "192.169.0.23"; + private static String elasticIp = "192.168.77.130"; private static int elasticPort = 9200; private static Logger logger = LoggerFactory.getLogger(EsHighLevelRestTest1.class); @@ -75,9 +75,9 @@ public static void main(String[] args) { queryById(); exists(); update(); -// deleteByQuery(); + deleteByQuery(); // deleteIndex(); -// delete(); + delete(); // bulk(); close(); } catch (Exception e) { @@ -88,10 +88,11 @@ public static void main(String[] args) { private static void insert() throws IOException { String index = "test1"; - String type = "_doc"; +// String type = "_doc"; // 唯一编号 - String id = "1"; - IndexRequest request = new IndexRequest(index, type, id); +// String id = "1"; +// IndexRequest request = new IndexRequest(index, type, id); + IndexRequest request = new IndexRequest(index); /* * 第一种方式,通过jsonString进行创建 */ @@ -186,7 +187,7 @@ private static void close() { private static void createIndex() throws IOException { // 类型 - String type = "_doc"; +// String type = "_doc"; String index = "test1"; // setting 的值 Map setmapping = new HashMap<>(); @@ -215,11 +216,11 @@ private static void createIndex() throws IOException { properties.put("sendtime", date); Map mapping = new HashMap<>(); mapping.put("properties", properties); - jsonMap2.put(type, mapping); +// jsonMap2.put(type, mapping); GetIndexRequest getRequest = new GetIndexRequest(); getRequest.indices(index); - getRequest.types(type); +// getRequest.types(type); getRequest.local(false); getRequest.humanReadable(true); boolean exists2 = client.indices().exists(getRequest, RequestOptions.DEFAULT); @@ -234,7 +235,7 @@ private static void createIndex() throws IOException { // 加载数据类型 request.settings(setmapping); //设置mapping参数 - request.mapping(type, jsonMap2); +// request.mapping(type, jsonMap2); //设置别名 request.alias(new Alias("pancm_alias")); CreateIndexResponse createIndexResponse = client.indices().create(request, RequestOptions.DEFAULT); @@ -254,12 +255,16 @@ private static void createIndex() throws IOException { * * @throws IOException */ - private static void deleteIndex() throws IOException { - String index = "userindex"; + private static void deleteIndex(){ + String index = "test1"; DeleteIndexRequest request = new DeleteIndexRequest(index); - // 同步删除 - client.indices().delete(request,RequestOptions.DEFAULT); - System.out.println("删除索引库成功!"+index); + try { + // 同步删除 + client.indices().delete(request,RequestOptions.DEFAULT); + System.out.println("删除索引库成功!"+index); + } catch (IOException e) { + throw new RuntimeException(e); + } } @@ -269,12 +274,14 @@ private static void deleteIndex() throws IOException { * @throws IOException */ private static void queryById() { - String type = "_doc"; +// String type = "_doc"; String index = "test1"; // 唯一编号 String id = "1"; // 创建查询请求 - GetRequest getRequest = new GetRequest(index, type, id); +// GetRequest getRequest = new GetRequest(index, type, id); + GetRequest getRequest = new GetRequest(index); + getRequest.id(id); GetResponse getResponse = null; try { @@ -314,6 +321,7 @@ private static void exists() throws IOException { String id = "1"; // 创建查询请求 GetRequest getRequest = new GetRequest(index, type, id); +// GetRequest getRequest = new GetRequest(index); boolean exists = client.exists(getRequest, RequestOptions.DEFAULT); @@ -341,14 +349,14 @@ public void onFailure(Exception e) { * @throws IOException */ private static void update() throws IOException { - String type = "_doc"; +// String type = "_doc"; String index = "test1"; // 唯一编号 String id = "1"; UpdateRequest upateRequest = new UpdateRequest(); upateRequest.id(id); upateRequest.index(index); - upateRequest.type(type); +// upateRequest.type(type); // 依旧可以使用Map这种集合作为更新条件 Map jsonMap = new HashMap<>(); @@ -373,10 +381,11 @@ private static void update() throws IOException { * @throws IOException */ private static void updateByQuery() throws IOException { - String type = "_doc"; +// String type = "_doc"; String index = "test1"; // - UpdateByQueryRequest request = new UpdateByQueryRequest(index,type); +// UpdateByQueryRequest request = new UpdateByQueryRequest(index,type); + UpdateByQueryRequest request = new UpdateByQueryRequest(index); // 设置查询条件 request.setQuery(new TermQueryBuilder("user", "pancm")); BoolQueryBuilder boolQueryBuilder = new BoolQueryBuilder(); @@ -427,10 +436,10 @@ private static void delete() throws IOException { String index = "test1"; // 唯一编号 String id = "1"; - DeleteRequest deleteRequest = new DeleteRequest(); + DeleteRequest deleteRequest = new DeleteRequest(index); deleteRequest.id(id); - deleteRequest.index(index); - deleteRequest.type(type); +// deleteRequest.index(index); +// deleteRequest.type(type); // 设置超时时间 deleteRequest.timeout(TimeValue.timeValueMinutes(2)); // 设置刷新策略"wait_for" @@ -481,11 +490,11 @@ public void onFailure(Exception e) { * @throws IOException */ private static void deleteByQuery() throws IOException { - String type = "_doc"; +// String type = "_doc"; String index = "test1"; - DeleteByQueryRequest request = new DeleteByQueryRequest(index,type); + DeleteByQueryRequest request = new DeleteByQueryRequest(index); // 设置查询条件 - request.setQuery(QueryBuilders.termQuery("uid",1234)); + request.setQuery(QueryBuilders.termQuery("uid",12345)); // 同步执行 BulkByScrollResponse bulkResponse = client.deleteByQuery(request, RequestOptions.DEFAULT); diff --git a/src/main/java/com/pancm/elasticsearch/EsHighLevelRestTest2.java b/src/main/java/com/pancm/elasticsearch/EsHighLevelRestTest2.java index 68a339a..9f59cb7 100644 --- a/src/main/java/com/pancm/elasticsearch/EsHighLevelRestTest2.java +++ b/src/main/java/com/pancm/elasticsearch/EsHighLevelRestTest2.java @@ -39,7 +39,7 @@ */ public class EsHighLevelRestTest2 { - private static String elasticIp = "192.169.0.23"; + private static String elasticIp = "192.168.77.130"; private static int elasticPort = 9200; private static Logger logger = LoggerFactory.getLogger(EsHighLevelRestTest2.class); diff --git a/src/main/java/com/pancm/elasticsearch/EsUtil.java b/src/main/java/com/pancm/elasticsearch/EsUtil.java index 31b2dcc..f6ac160 100644 --- a/src/main/java/com/pancm/elasticsearch/EsUtil.java +++ b/src/main/java/com/pancm/elasticsearch/EsUtil.java @@ -65,7 +65,7 @@ public static void main(String[] args) { try { - EsUtil.build("192.169.0.23:9200"); + EsUtil.build("192.168.77.130:9200"); System.out.println("ES连接初始化成功!"); // createIndexTest(); // System.out.println("ES索引库创建成功!"); diff --git a/src/main/java/com/pancm/jdk8/StreamTest.java b/src/main/java/com/pancm/jdk8/StreamTest.java index 7108409..1a34d17 100644 --- a/src/main/java/com/pancm/jdk8/StreamTest.java +++ b/src/main/java/com/pancm/jdk8/StreamTest.java @@ -113,18 +113,23 @@ private static void test2() { stream = list.stream(); /* - * 流之间的相互转化 一个 Stream 只可以使用一次,这段代码为了简洁而重复使用了数次,因此会抛出异常 + * 流之间的相互转化 一个 Stream 只可以使用一次,这段代码为了简洁而重复使用了数次,因此会抛出异常.已修复。 */ try { + //同一个流 连使用两次会报错:java.lang.IllegalStateException: stream has already been operated upon or closed + Stream stream1 = Stream.of("a", "b", "c"); Stream stream2 = Stream.of("a", "b", "c"); + Stream stream3 = Stream.of("a", "b", "c"); + Stream stream4 = Stream.of("a", "b", "c"); + Stream stream5 = Stream.of("a", "b", "c"); // 转换成 Array - String[] strArray1 = stream2.toArray(String[]::new); + String[] strArray1 = stream1.toArray(String[]::new); // 转换成 Collection List list1 = stream2.collect(Collectors.toList()); - List list2 = stream2.collect(Collectors.toCollection(ArrayList::new)); - Set set1 = stream2.collect(Collectors.toSet()); - Stack stack1 = stream2.collect(Collectors.toCollection(Stack::new)); + List list2 = stream3.collect(Collectors.toCollection(ArrayList::new)); + Set set1 = stream4.collect(Collectors.toSet()); + Stack stack1 = stream5.collect(Collectors.toCollection(Stack::new)); // 转换成 String String str = stream.collect(Collectors.joining()).toString(); @@ -257,7 +262,7 @@ private static void test2() { User user = new User(i, "pancm" + i); list9.add(user); } - System.out.println("截取之前的数据:"); + System.out.println("截取之前的数据:");//在User对象中.每次调用getName打印 // 取前3条数据,但是扔掉了前面的2条,可以理解为拿到的数据为 2<=i<3 (i 是数值下标) List list10 = list9.stream().map(User::getName).limit(3).skip(2).collect(Collectors.toList()); System.out.println("截取之后的数据:" + list10); diff --git a/src/main/java/com/pancm/mq/kafka/examples/WordCountDemo.java b/src/main/java/com/pancm/mq/kafka/examples/WordCountDemo.java index 7298806..d4012b9 100644 --- a/src/main/java/com/pancm/mq/kafka/examples/WordCountDemo.java +++ b/src/main/java/com/pancm/mq/kafka/examples/WordCountDemo.java @@ -33,7 +33,7 @@ public class WordCountDemo { public static void main(String[] args) throws Exception { Properties props = new Properties(); props.put(StreamsConfig.APPLICATION_ID_CONFIG, "streams-wordcount"); - props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "192.169.0.23:9092"); + props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "192.168.77.130:9092"); props.put(StreamsConfig.CACHE_MAX_BYTES_BUFFERING_CONFIG, 0); props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName()); props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName()); diff --git a/src/main/java/com/pancm/mq/kafka/others/TestConsumer.java b/src/main/java/com/pancm/mq/kafka/others/TestConsumer.java index b1b9d93..3bed765 100644 --- a/src/main/java/com/pancm/mq/kafka/others/TestConsumer.java +++ b/src/main/java/com/pancm/mq/kafka/others/TestConsumer.java @@ -12,7 +12,7 @@ public class TestConsumer { public static void main(String[] args) { Properties props = new Properties(); - props.put("bootstrap.servers", "192.169.0.23:9092"); + props.put("bootstrap.servers", "192.168.77.130:9092"); System.out.println("this is the group part test 1"); //消费者的组id props.put("group.id", "GroupA");//这里是GroupA或者GroupB diff --git a/src/main/java/com/pancm/mq/kafka/others/TestProducer.java b/src/main/java/com/pancm/mq/kafka/others/TestProducer.java index 88db1a7..f50ca94 100644 --- a/src/main/java/com/pancm/mq/kafka/others/TestProducer.java +++ b/src/main/java/com/pancm/mq/kafka/others/TestProducer.java @@ -10,7 +10,7 @@ public class TestProducer { public static void main(String[] args) { System.out.println("开始..."); Properties props = new Properties(); - props.put("bootstrap.servers", "192.169.0.23:9092"); + props.put("bootstrap.servers", "192.168.77.130:9092"); //The "all" setting we have specified will result in blocking on the full commit of the record, the slowest but most durable setting. //“所有”设置将导致记录的完整提交阻塞,最慢的,但最持久的设置。 props.put("acks", "all"); diff --git a/src/main/java/com/pancm/mq/kafka/test1/KafkaConsumerTest.java b/src/main/java/com/pancm/mq/kafka/test1/KafkaConsumerTest.java index 4b5a188..2495922 100644 --- a/src/main/java/com/pancm/mq/kafka/test1/KafkaConsumerTest.java +++ b/src/main/java/com/pancm/mq/kafka/test1/KafkaConsumerTest.java @@ -13,7 +13,7 @@ * * Title: KafkaConsumerTest * Description: -* kafka消费者 demo +* kafka消费者 demo:https://www.cnblogs.com/xuwujing/p/8371127.html * Version:1.0.0 * @author pancm * @date 2018年1月26日 @@ -30,7 +30,8 @@ public class KafkaConsumerTest implements Runnable { public KafkaConsumerTest(String topicName) { Properties props = new Properties(); //kafka消费的的地址 - props.put("bootstrap.servers", "master:9092,slave1:9092,slave2:9092"); +// props.put("bootstrap.servers", "master:9092,slave1:9092,slave2:9092"); + props.put("bootstrap.servers", "localhost:9092"); //组名 不同组名可以重复消费 props.put("group.id", GROUPID); //是否自动提交 diff --git a/src/main/java/com/pancm/mq/kafka/test1/KafkaProducerTest.java b/src/main/java/com/pancm/mq/kafka/test1/KafkaProducerTest.java index 728941a..9f8c677 100644 --- a/src/main/java/com/pancm/mq/kafka/test1/KafkaProducerTest.java +++ b/src/main/java/com/pancm/mq/kafka/test1/KafkaProducerTest.java @@ -24,7 +24,7 @@ public class KafkaProducerTest implements Runnable { public KafkaProducerTest(String topicName) { Properties props = new Properties(); - props.put("bootstrap.servers", "master:9092,slave1:9092,slave2:9092"); + props.put("bootstrap.servers", "localhost:9092"); //acks=0:如果设置为0,生产者不会等待kafka的响应。 //acks=1:这个配置意味着kafka会把这条消息写到本地日志文件中,但是不会等待集群中其他机器的成功响应。 //acks=all:这个配置意味着leader会等待所有的follower同步完成。这个确保消息不会丢失,除非kafka集群中所有机器挂掉。这是最强的可用性保证。 diff --git a/src/main/java/com/pancm/mq/kafka/test2/KafkaConsumerTest.java b/src/main/java/com/pancm/mq/kafka/test2/KafkaConsumerTest.java index b250196..89840fe 100644 --- a/src/main/java/com/pancm/mq/kafka/test2/KafkaConsumerTest.java +++ b/src/main/java/com/pancm/mq/kafka/test2/KafkaConsumerTest.java @@ -16,7 +16,7 @@ * Title: KafkaConsumerTest * Description: * kafka消费者 demo -* 手动提交测试 +* 手动提交测试:https://www.cnblogs.com/xuwujing/p/8432984.html * Version:1.0.0 * @author pancm * @date 2018年1月26日 @@ -77,7 +77,7 @@ public void run() { private void init() { Properties props = new Properties(); //kafka消费的的地址 - props.put("bootstrap.servers", "master:9092,slave1:9092,slave2:9092"); + props.put("bootstrap.servers", "localhost:9092"); //组名 不同组名可以重复消费 props.put("group.id", GROUPID); //是否自动提交 diff --git a/src/main/java/com/pancm/mq/kafka/test2/KafkaProducerTest.java b/src/main/java/com/pancm/mq/kafka/test2/KafkaProducerTest.java index e1b11f2..e271aff 100644 --- a/src/main/java/com/pancm/mq/kafka/test2/KafkaProducerTest.java +++ b/src/main/java/com/pancm/mq/kafka/test2/KafkaProducerTest.java @@ -24,7 +24,7 @@ public class KafkaProducerTest implements Runnable { public KafkaProducerTest(String topicName) { Properties props = new Properties(); - props.put("bootstrap.servers", "master:9092,slave1:9092,slave2:9092"); + props.put("bootstrap.servers", "localhost:9092"); //acks=0:如果设置为0,生产者不会等待kafka的响应。 //acks=1:这个配置意味着kafka会把这条消息写到本地日志文件中,但是不会等待集群中其他机器的成功响应。 //acks=all:这个配置意味着leader会等待所有的follower同步完成。这个确保消息不会丢失,除非kafka集群中所有机器挂掉。这是最强的可用性保证。 diff --git a/src/main/java/com/pancm/mq/kafka/test3/KafkaConsumerTest.java b/src/main/java/com/pancm/mq/kafka/test3/KafkaConsumerTest.java index 1d73761..607dd3d 100644 --- a/src/main/java/com/pancm/mq/kafka/test3/KafkaConsumerTest.java +++ b/src/main/java/com/pancm/mq/kafka/test3/KafkaConsumerTest.java @@ -23,8 +23,8 @@ public class KafkaConsumerTest extends Thread { private ConsumerRecords msgList; private final String topic; private static final String GROUPID = "groupA1"; - private final String servers="master:9092,slave1:9092,slave2:9092"; - + private final String servers="localhost:9092"; + public KafkaConsumerTest(String topicName) { Properties props = new Properties(); //kafka消费的的地址 diff --git a/src/main/java/com/pancm/mq/kafka/test3/KafkaConsumerTest3.java b/src/main/java/com/pancm/mq/kafka/test3/KafkaConsumerTest3.java index 60d8e34..29b3192 100644 --- a/src/main/java/com/pancm/mq/kafka/test3/KafkaConsumerTest3.java +++ b/src/main/java/com/pancm/mq/kafka/test3/KafkaConsumerTest3.java @@ -93,7 +93,7 @@ private long getOffset(int partId) { private void init() { Properties props = new Properties(); //kafka消费的的地址 - props.put("bootstrap.servers", "master:9092,slave1:9092,slave2:9092"); + props.put("bootstrap.servers", "localhost:9092"); //组名 不同组名可以重复消费 props.put("group.id", GROUPID); //是否自动提交 diff --git a/src/main/java/com/pancm/mq/kafka/test3/KafkaProducerTest.java b/src/main/java/com/pancm/mq/kafka/test3/KafkaProducerTest.java index c206571..5c3663f 100644 --- a/src/main/java/com/pancm/mq/kafka/test3/KafkaProducerTest.java +++ b/src/main/java/com/pancm/mq/kafka/test3/KafkaProducerTest.java @@ -21,7 +21,7 @@ public class KafkaProducerTest implements Runnable { private final KafkaProducer producer; private final String topic; private int k=10; - private final String servers="master:9092,slave1:9092,slave2:9092"; + private final String servers="localhost:9092"; /** * @param topic 消息名称 * @param diff --git a/src/main/java/com/pancm/mq/rocketmq/Consumer.java b/src/main/java/com/pancm/mq/rocketmq/Consumer.java new file mode 100644 index 0000000..882f742 --- /dev/null +++ b/src/main/java/com/pancm/mq/rocketmq/Consumer.java @@ -0,0 +1,41 @@ +package com.pancm.mq.rocketmq; + +import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer; +import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext; +import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus; +import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently; +import org.apache.rocketmq.client.exception.MQClientException; +import org.apache.rocketmq.common.message.MessageExt; + +import java.util.List; + +public class Consumer { + public static void main(String[] args) { + DefaultMQPushConsumer mqConsumer = new + DefaultMQPushConsumer("consumer_group"); + mqConsumer.setNamesrvAddr("localhost:9876"); + + // 设置消息监听器 + mqConsumer.registerMessageListener(new MessageListenerConcurrently() { + @Override + public ConsumeConcurrentlyStatus consumeMessage(List list, ConsumeConcurrentlyContext consumeConcurrentlyContext) { + MessageExt messageExt = list.get(0); + byte[] body = messageExt.getBody(); + System.out.printf(String.valueOf(body.length)); + System.out.printf(messageExt.getMsgId()); + return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; + } + }); + //订阅主题 + try { + mqConsumer.subscribe("topicA", "*"); + } catch (MQClientException e) { + throw new RuntimeException(e); + } + try { + mqConsumer.start(); + } catch (MQClientException e) { + throw new RuntimeException(e); + } + } +} \ No newline at end of file diff --git a/src/main/java/com/pancm/mq/rocketmq/Producer.java b/src/main/java/com/pancm/mq/rocketmq/Producer.java new file mode 100644 index 0000000..f91e9cf --- /dev/null +++ b/src/main/java/com/pancm/mq/rocketmq/Producer.java @@ -0,0 +1,31 @@ +package com.pancm.mq.rocketmq; + +import org.apache.rocketmq.client.exception.MQBrokerException; +import org.apache.rocketmq.client.exception.MQClientException; +import org.apache.rocketmq.client.producer.DefaultMQProducer; +import org.apache.rocketmq.client.producer.SendResult; +import org.apache.rocketmq.common.message.Message; +import org.apache.rocketmq.remoting.exception.RemotingException; + +import java.nio.charset.StandardCharsets; + +public class Producer { + public static void main(String[] args) { + DefaultMQProducer producer = new DefaultMQProducer("producer_group"); + producer.setNamesrvAddr("127.0.0.1:9876"); + try { + producer.start(); + } catch (MQClientException e) { + throw new RuntimeException(e); + } + //待发送消息 + Message topicA = new Message("topicA", "message".getBytes(StandardCharsets.UTF_8)); + try { + SendResult sendResult = producer.send(topicA); + System.out.printf(sendResult.toString()); + } catch (MQClientException | RemotingException | MQBrokerException | InterruptedException e) { + throw new RuntimeException(e); + } + + } +} diff --git a/src/main/java/com/pancm/others/LombokTest.java b/src/main/java/com/pancm/others/LombokTest.java index 3072684..afb5b80 100644 --- a/src/main/java/com/pancm/others/LombokTest.java +++ b/src/main/java/com/pancm/others/LombokTest.java @@ -101,8 +101,8 @@ public static void main(String[] args) { */ private static void test1() throws IOException { //创建要操作的文件路径和名称 - String path ="E:/test/hello.txt"; - String str="你好!"; + String path = "d:/hello.txt"; + String str = "你好!"; FileWriter fw = new FileWriter(path); BufferedWriter bw=new BufferedWriter(fw); bw.write(str); @@ -118,8 +118,8 @@ private static void test1() throws IOException { */ private static void test2() throws IOException { //创建要操作的文件路径和名称 - String path ="E:/test/hello.txt"; - String str="你好!"; + String path = "d:/hello.txt"; + String str = "你好!"; @Cleanup FileWriter fw = new FileWriter(path); @Cleanup @@ -132,7 +132,6 @@ private static void test2() throws IOException { /** - * @param name2 */ private static void test3(String name) { if(null == name) { @@ -142,7 +141,6 @@ private static void test3(String name) { } /** - * @param name2 */ private static void test4(@NonNull String name) { System.out.println(name); diff --git a/src/main/java/com/pancm/others/body.json b/src/main/java/com/pancm/others/body.json new file mode 100644 index 0000000..bee0294 --- /dev/null +++ b/src/main/java/com/pancm/others/body.json @@ -0,0 +1,7 @@ +{ + "appid": "Sg8", + "ukey": "03e1923c", + "authType": "3", + "time": "{{time}}", + "sign": "{{sign}}" +} \ No newline at end of file diff --git a/src/main/java/com/pancm/others/getToken.java b/src/main/java/com/pancm/others/getToken.java new file mode 100644 index 0000000..457772e --- /dev/null +++ b/src/main/java/com/pancm/others/getToken.java @@ -0,0 +1,84 @@ +package com.pancm.others; + +import com.alibaba.fastjson.JSON; +import org.apache.commons.codec.digest.DigestUtils; +import org.apache.http.HttpResponse; +import org.apache.http.client.methods.HttpPost; +import org.apache.http.entity.ContentType; +import org.apache.http.entity.StringEntity; +import org.apache.http.impl.client.CloseableHttpClient; +import org.apache.http.impl.client.HttpClients; +import org.apache.http.util.EntityUtils; + +import java.nio.charset.StandardCharsets; +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.StringJoiner; + +public class getToken { + public static void main(String[] args) { + /** + * + * postman 脚本 + * 使用pre-request-script.js + * 请求体:body.json + * {{time}}:获取环境变量值。json请求中需要加双引号"" + * + */ + + /** + * java 代码 + */ + + String baseUrl = "baseurl"; + String url = baseUrl + "/auth/token/get"; + try (CloseableHttpClient httpClient = HttpClients.createDefault()) { + String ukey = "ukeyStr"; + String usecret = "usecret"; + Map queryMap = new HashMap<>(); + queryMap.put("ukey", ukey); + queryMap.put("appid", "appid"); + queryMap.put("authType", "3"); + queryMap.put("time", (System.currentTimeMillis()) + ""); + String bodyStr = ukey + createLinkString(queryMap, null) + usecret; + System.out.println(bodyStr); + + String sign = DigestUtils.md5Hex(bodyStr); + HttpPost method = new HttpPost(url); + + method.addHeader("WR-Client-Id", "12345"); + method.addHeader("WR-ukey", ukey); + queryMap.put("sign", sign); + System.out.println(JSON.toJSONString(queryMap)); + StringEntity entity = new StringEntity(JSON.toJSONString(queryMap), ContentType.APPLICATION_JSON); + method.setEntity(entity); + + HttpResponse response = httpClient.execute(method); + System.out.println(response.getStatusLine().getStatusCode()); + System.out.println(EntityUtils.toString(response.getEntity(), StandardCharsets.UTF_8)); + } catch (Exception e) { + e.printStackTrace(); + } + } + + private static String createLinkString(Map params, String joinStr) { + if (params == null || params.isEmpty()) { + return null; + } + List keys = new ArrayList<>(params.keySet()); + Collections.sort(keys); + //默认用"&"连接 + StringJoiner joiner = new StringJoiner((joinStr == null) ? "&" : joinStr); + keys.forEach(key -> { + try { + joiner.add(key + "=" + params.get(key)); + } catch (Exception e) { + return; + } + }); + return joiner.toString(); + } +} diff --git a/src/main/java/com/pancm/others/pre-request-script.js b/src/main/java/com/pancm/others/pre-request-script.js new file mode 100644 index 0000000..41557d2 --- /dev/null +++ b/src/main/java/com/pancm/others/pre-request-script.js @@ -0,0 +1,22 @@ +var appid = "Sg8"; +var ukey = "03e1923c"; +var authType = "3"; +var usecret = "442"; +var time = (new Date()).getTime().toString();//获取当前时间戳 +pm.environment.set('time', time); + +var str = 'appid=' + appid + '&authType=' + authType + '&time=' + time + '&ukey=' + ukey; + +var bodyStr = ukey + str + usecret; +var sign = CryptoJS.MD5(bodyStr).toString(); +pm.environment.set('sign', sign);//加密后的密码字符串赋值给环境变量 + + +//获取body中参数,上面代码可以进行改造 +var body = pm.request.body.raw +var body_json = JSON.parse(body) +pwd = body_json["password"] +console.log(pwd) //在console打印pwd参数 +//获取header中的参数 +pm.request.headers.get("Cookie") +