2 * ============LICENSE_START=======================================================
4 * ================================================================================
5 * Copyright 2019 China Mobile
6 *=================================================================================
7 * Licensed under the Apache License, Version 2.0 (the "License");
8 * you may not use this file except in compliance with the License.
9 * You may obtain a copy of the License at
11 * http://www.apache.org/licenses/LICENSE-2.0
13 * Unless required by applicable law or agreed to in writing, software
14 * distributed under the License is distributed on an "AS IS" BASIS,
15 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
16 * See the License for the specific language governing permissions and
17 * limitations under the License.
18 * ============LICENSE_END=========================================================
21 package org.onap.datalake.feeder.service;
23 import java.util.ArrayList;
24 import java.util.List;
25 import java.util.UUID;
27 import javax.annotation.PostConstruct;
28 import javax.annotation.PreDestroy;
30 import org.json.JSONObject;
31 import org.onap.datalake.feeder.config.ApplicationConfiguration;
32 import org.onap.datalake.feeder.domain.Db;
33 import org.onap.datalake.feeder.dto.TopicConfig;
34 import org.slf4j.Logger;
35 import org.slf4j.LoggerFactory;
36 import org.springframework.beans.factory.annotation.Autowired;
37 import org.springframework.stereotype.Service;
39 import com.couchbase.client.java.Bucket;
40 import com.couchbase.client.java.Cluster;
41 import com.couchbase.client.java.CouchbaseCluster;
42 import com.couchbase.client.java.document.JsonDocument;
43 import com.couchbase.client.java.document.json.JsonObject;
44 import com.couchbase.client.java.env.CouchbaseEnvironment;
45 import com.couchbase.client.java.env.DefaultCouchbaseEnvironment;
46 import com.couchbase.client.java.error.DocumentAlreadyExistsException;
49 import rx.functions.Func1;
52 * Service to use Couchbase
58 public class CouchbaseService {
60 private final Logger log = LoggerFactory.getLogger(this.getClass());
63 ApplicationConfiguration config;
66 private DbService dbService;
69 private boolean isReady = false;
73 // Initialize Couchbase Connection
75 Db couchbase = dbService.getCouchbase();
77 //this tunes the SDK (to customize connection timeout)
78 CouchbaseEnvironment env = DefaultCouchbaseEnvironment.builder().connectTimeout(60000) // 60s, default is 5s
80 Cluster cluster = CouchbaseCluster.create(env, couchbase.getHost());
81 cluster.authenticate(couchbase.getLogin(), couchbase.getPass());
82 bucket = cluster.openBucket(couchbase.getDatabase());
83 // Create a N1QL Primary Index (but ignore if it exists)
84 bucket.bucketManager().createN1qlPrimaryIndex(true, false);
86 log.info("Connected to Couchbase {} as {}", couchbase.getHost(), couchbase.getLogin());
88 } catch (Exception ex) {
89 log.error("error connection to Couchbase.", ex);
95 public void cleanUp() {
96 config.getShutdownLock().readLock().lock();
99 log.info("bucket.close() at cleanUp.");
102 config.getShutdownLock().readLock().unlock();
106 public void saveJsons(TopicConfig topic, List<JSONObject> jsons) {
107 List<JsonDocument> documents = new ArrayList<>(jsons.size());
108 for (JSONObject json : jsons) {
109 //convert to Couchbase JsonObject from org.json JSONObject
110 JsonObject jsonObject = JsonObject.fromJson(json.toString());
112 long timestamp = jsonObject.getLong(config.getTimestampLabel());//this is Kafka time stamp, which is added in StoreService.messageToJson()
115 int expiry = (int) (timestamp / 1000L) + topic.getTtl() * 3600 * 24; //in second
117 String id = getId(topic, json);
118 JsonDocument doc = JsonDocument.create(id, expiry, jsonObject);
122 saveDocuments(documents);
123 } catch (DocumentAlreadyExistsException e) {
124 log.error("Some or all the following ids are duplicate.");
125 for(JsonDocument document : documents) {
126 log.error("saveJsons() DocumentAlreadyExistsException {}", document.id());
128 } catch (rx.exceptions.CompositeException e) {
129 List<Throwable> causes = e.getExceptions();
130 for(Throwable cause : causes) {
131 log.error("saveJsons() CompositeException cause {}", cause.getMessage());
133 } catch (Exception e) {
134 log.error("error saving to Couchbase.", e);
136 log.debug("saved text to topic = {}, this batch count = {} ", topic, documents.size());
139 public String getId(TopicConfig topic, JSONObject json) {
140 //if this topic requires extract id from JSON
141 String id = topic.getMessageId(json);
146 String topicStr = topic.getName();
147 id = topicStr+":"+UUID.randomUUID();
149 //https://forums.couchbase.com/t/how-to-set-an-auto-increment-id/4892/2
150 //atomically get the next sequence number:
151 // increment by 1, initialize at 0 if counter doc not found
152 //TODO how slow is this compared with above UUID approach?
153 //sometimes this gives java.util.concurrent.TimeoutException
154 //JsonLongDocument nextIdNumber = bucket.counter(topicStr, 1, 0); //like 12345
155 //id = topicStr + ":" + nextIdNumber.content();
160 //https://docs.couchbase.com/java-sdk/2.7/document-operations.html
161 private void saveDocuments(List<JsonDocument> documents) {
162 Observable.from(documents).flatMap(new Func1<JsonDocument, Observable<JsonDocument>>() {
164 public Observable<JsonDocument> call(final JsonDocument docToInsert) {
165 return bucket.async().insert(docToInsert);
167 }).last().toBlocking().single();