2021-09-05 17:10:14 +00:00
|
|
|
"""
|
|
|
|
Functionality:
|
|
|
|
- initial elastic search setup
|
|
|
|
- index configuration is represented in INDEX_CONFIG
|
|
|
|
- index mapping and settings validation
|
|
|
|
- backup and restore
|
|
|
|
"""
|
|
|
|
|
|
|
|
import json
|
|
|
|
import os
|
2021-09-16 10:34:20 +00:00
|
|
|
import zipfile
|
2021-09-05 17:10:14 +00:00
|
|
|
from datetime import datetime
|
|
|
|
|
|
|
|
import requests
|
|
|
|
from home.src.config import AppConfig
|
2021-09-25 11:59:54 +00:00
|
|
|
from home.src.helper import ignore_filelist
|
2021-09-05 17:10:14 +00:00
|
|
|
|
|
|
|
# expected mapping and settings
|
|
|
|
INDEX_CONFIG = [
|
|
|
|
{
|
2021-09-21 09:25:22 +00:00
|
|
|
"index_name": "channel",
|
|
|
|
"expected_map": {
|
2021-09-05 17:10:14 +00:00
|
|
|
"channel_id": {
|
|
|
|
"type": "keyword",
|
|
|
|
},
|
|
|
|
"channel_name": {
|
|
|
|
"type": "text",
|
|
|
|
"fields": {
|
|
|
|
"keyword": {
|
|
|
|
"type": "keyword",
|
|
|
|
"ignore_above": 256,
|
2021-09-21 09:25:22 +00:00
|
|
|
"normalizer": "to_lower",
|
2021-09-05 17:10:14 +00:00
|
|
|
},
|
|
|
|
"search_as_you_type": {
|
|
|
|
"type": "search_as_you_type",
|
|
|
|
"doc_values": False,
|
2021-09-21 09:25:22 +00:00
|
|
|
"max_shingle_size": 3,
|
|
|
|
},
|
|
|
|
},
|
2021-09-05 17:10:14 +00:00
|
|
|
},
|
2021-09-21 09:25:22 +00:00
|
|
|
"channel_banner_url": {"type": "keyword", "index": False},
|
|
|
|
"channel_thumb_url": {"type": "keyword", "index": False},
|
|
|
|
"channel_description": {"type": "text"},
|
|
|
|
"channel_last_refresh": {"type": "date", "format": "epoch_second"},
|
2021-09-05 17:10:14 +00:00
|
|
|
},
|
2021-09-21 09:25:22 +00:00
|
|
|
"expected_set": {
|
2021-09-05 17:10:14 +00:00
|
|
|
"analysis": {
|
|
|
|
"normalizer": {
|
2021-09-21 09:25:22 +00:00
|
|
|
"to_lower": {"type": "custom", "filter": ["lowercase"]}
|
2021-09-05 17:10:14 +00:00
|
|
|
}
|
|
|
|
},
|
2021-09-21 09:25:22 +00:00
|
|
|
"number_of_replicas": "0",
|
|
|
|
},
|
2021-09-05 17:10:14 +00:00
|
|
|
},
|
|
|
|
{
|
2021-09-21 09:25:22 +00:00
|
|
|
"index_name": "video",
|
|
|
|
"expected_map": {
|
|
|
|
"vid_thumb_url": {"type": "text", "index": False},
|
|
|
|
"date_downloaded": {"type": "date"},
|
2021-09-05 17:10:14 +00:00
|
|
|
"channel": {
|
|
|
|
"properties": {
|
|
|
|
"channel_id": {
|
|
|
|
"type": "keyword",
|
|
|
|
},
|
|
|
|
"channel_name": {
|
|
|
|
"type": "text",
|
|
|
|
"fields": {
|
|
|
|
"keyword": {
|
|
|
|
"type": "keyword",
|
|
|
|
"ignore_above": 256,
|
2021-09-21 09:25:22 +00:00
|
|
|
"normalizer": "to_lower",
|
2021-09-05 17:10:14 +00:00
|
|
|
},
|
|
|
|
"search_as_you_type": {
|
|
|
|
"type": "search_as_you_type",
|
|
|
|
"doc_values": False,
|
2021-09-21 09:25:22 +00:00
|
|
|
"max_shingle_size": 3,
|
|
|
|
},
|
|
|
|
},
|
2021-09-05 17:10:14 +00:00
|
|
|
},
|
2021-09-21 09:25:22 +00:00
|
|
|
"channel_banner_url": {"type": "keyword", "index": False},
|
|
|
|
"channel_thumb_url": {"type": "keyword", "index": False},
|
|
|
|
"channel_description": {"type": "text"},
|
2021-09-05 17:10:14 +00:00
|
|
|
"channel_last_refresh": {
|
|
|
|
"type": "date",
|
2021-09-21 09:25:22 +00:00
|
|
|
"format": "epoch_second",
|
|
|
|
},
|
2021-09-05 17:10:14 +00:00
|
|
|
}
|
|
|
|
},
|
2021-09-21 09:25:22 +00:00
|
|
|
"description": {"type": "text"},
|
|
|
|
"media_url": {"type": "keyword", "index": False},
|
2021-09-05 17:10:14 +00:00
|
|
|
"title": {
|
|
|
|
"type": "text",
|
|
|
|
"fields": {
|
|
|
|
"keyword": {
|
|
|
|
"type": "keyword",
|
|
|
|
"ignore_above": 256,
|
2021-09-21 09:25:22 +00:00
|
|
|
"normalizer": "to_lower",
|
2021-09-05 17:10:14 +00:00
|
|
|
},
|
|
|
|
"search_as_you_type": {
|
|
|
|
"type": "search_as_you_type",
|
|
|
|
"doc_values": False,
|
2021-09-21 09:25:22 +00:00
|
|
|
"max_shingle_size": 3,
|
|
|
|
},
|
|
|
|
},
|
2021-09-05 17:10:14 +00:00
|
|
|
},
|
2021-09-21 09:25:22 +00:00
|
|
|
"vid_last_refresh": {"type": "date"},
|
|
|
|
"youtube_id": {"type": "keyword"},
|
|
|
|
"published": {"type": "date"},
|
2021-11-11 10:56:29 +00:00
|
|
|
"playlist": {
|
2021-11-13 10:34:58 +00:00
|
|
|
"type": "text",
|
|
|
|
"fields": {
|
|
|
|
"keyword": {
|
|
|
|
"type": "keyword",
|
|
|
|
"ignore_above": 256,
|
|
|
|
"normalizer": "to_lower",
|
|
|
|
}
|
|
|
|
},
|
2021-11-11 10:56:29 +00:00
|
|
|
},
|
2021-09-05 17:10:14 +00:00
|
|
|
},
|
2021-09-21 09:25:22 +00:00
|
|
|
"expected_set": {
|
2021-09-05 17:10:14 +00:00
|
|
|
"analysis": {
|
|
|
|
"normalizer": {
|
2021-09-21 09:25:22 +00:00
|
|
|
"to_lower": {"type": "custom", "filter": ["lowercase"]}
|
2021-09-05 17:10:14 +00:00
|
|
|
}
|
|
|
|
},
|
2021-09-21 09:25:22 +00:00
|
|
|
"number_of_replicas": "0",
|
|
|
|
},
|
2021-09-05 17:10:14 +00:00
|
|
|
},
|
|
|
|
{
|
2021-09-21 09:25:22 +00:00
|
|
|
"index_name": "download",
|
|
|
|
"expected_map": {
|
|
|
|
"timestamp": {"type": "date"},
|
|
|
|
"channel_id": {"type": "keyword"},
|
2021-09-05 17:10:14 +00:00
|
|
|
"channel_name": {
|
|
|
|
"type": "text",
|
|
|
|
"fields": {
|
|
|
|
"keyword": {
|
|
|
|
"type": "keyword",
|
|
|
|
"ignore_above": 256,
|
2021-09-21 09:25:22 +00:00
|
|
|
"normalizer": "to_lower",
|
2021-09-05 17:10:14 +00:00
|
|
|
}
|
2021-09-21 09:25:22 +00:00
|
|
|
},
|
2021-09-05 17:10:14 +00:00
|
|
|
},
|
2021-09-21 09:25:22 +00:00
|
|
|
"status": {"type": "keyword"},
|
2021-09-05 17:10:14 +00:00
|
|
|
"title": {
|
|
|
|
"type": "text",
|
|
|
|
"fields": {
|
|
|
|
"keyword": {
|
|
|
|
"type": "keyword",
|
|
|
|
"ignore_above": 256,
|
2021-09-21 09:25:22 +00:00
|
|
|
"normalizer": "to_lower",
|
2021-09-05 17:10:14 +00:00
|
|
|
}
|
2021-09-21 09:25:22 +00:00
|
|
|
},
|
2021-09-05 17:10:14 +00:00
|
|
|
},
|
2021-09-21 09:25:22 +00:00
|
|
|
"vid_thumb_url": {"type": "keyword"},
|
|
|
|
"youtube_id": {"type": "keyword"},
|
2021-09-05 17:10:14 +00:00
|
|
|
},
|
2021-09-21 09:25:22 +00:00
|
|
|
"expected_set": {
|
2021-09-05 17:10:14 +00:00
|
|
|
"analysis": {
|
|
|
|
"normalizer": {
|
2021-09-21 09:25:22 +00:00
|
|
|
"to_lower": {"type": "custom", "filter": ["lowercase"]}
|
2021-09-05 17:10:14 +00:00
|
|
|
}
|
|
|
|
},
|
2021-09-21 09:25:22 +00:00
|
|
|
"number_of_replicas": "0",
|
|
|
|
},
|
|
|
|
},
|
2021-11-08 07:54:17 +00:00
|
|
|
{
|
|
|
|
"index_name": "playlist",
|
|
|
|
"expected_map": {
|
2021-11-11 13:57:28 +00:00
|
|
|
"playlist_id": {"type": "keyword"},
|
2021-11-08 07:54:17 +00:00
|
|
|
"playlist_description": {"type": "text"},
|
2021-11-11 13:57:28 +00:00
|
|
|
"playlist_name": {
|
|
|
|
"type": "text",
|
|
|
|
"fields": {
|
|
|
|
"keyword": {
|
|
|
|
"type": "keyword",
|
|
|
|
"ignore_above": 256,
|
|
|
|
"normalizer": "to_lower",
|
|
|
|
}
|
|
|
|
},
|
|
|
|
},
|
|
|
|
"playlist_channel": {
|
|
|
|
"type": "text",
|
|
|
|
"fields": {
|
|
|
|
"keyword": {
|
|
|
|
"type": "keyword",
|
|
|
|
"ignore_above": 256,
|
|
|
|
"normalizer": "to_lower",
|
|
|
|
}
|
|
|
|
},
|
|
|
|
},
|
|
|
|
"playlist_channel_id": {"type": "keyword"},
|
|
|
|
"playlist_thumbnail": {"type": "keyword"},
|
2021-12-06 09:56:56 +00:00
|
|
|
"playlist_last_refresh": {"type": "date"},
|
2021-11-08 07:54:17 +00:00
|
|
|
},
|
|
|
|
"expected_set": {
|
2021-11-11 13:57:28 +00:00
|
|
|
"analysis": {
|
|
|
|
"normalizer": {
|
|
|
|
"to_lower": {"type": "custom", "filter": ["lowercase"]}
|
|
|
|
}
|
|
|
|
},
|
2021-11-08 07:54:17 +00:00
|
|
|
"number_of_replicas": "0",
|
|
|
|
},
|
|
|
|
},
|
2021-09-05 17:10:14 +00:00
|
|
|
]
|
|
|
|
|
|
|
|
|
|
|
|
class ElasticIndex:
|
|
|
|
"""
|
|
|
|
handle mapping and settings on elastic search for a given index
|
|
|
|
"""
|
|
|
|
|
|
|
|
CONFIG = AppConfig().config
|
2021-09-21 09:25:22 +00:00
|
|
|
ES_URL = CONFIG["application"]["es_url"]
|
2021-10-28 08:49:58 +00:00
|
|
|
ES_AUTH = CONFIG["application"]["es_auth"]
|
2021-09-21 09:25:22 +00:00
|
|
|
HEADERS = {"Content-type": "application/json"}
|
2021-09-05 17:10:14 +00:00
|
|
|
|
|
|
|
def __init__(self, index_name, expected_map, expected_set):
|
|
|
|
self.index_name = index_name
|
|
|
|
self.expected_map = expected_map
|
|
|
|
self.expected_set = expected_set
|
|
|
|
self.exists, self.details = self.index_exists()
|
|
|
|
|
|
|
|
def index_exists(self):
|
2021-09-21 09:25:22 +00:00
|
|
|
"""check if index already exists and return mapping if it does"""
|
2021-09-05 17:10:14 +00:00
|
|
|
index_name = self.index_name
|
2021-09-21 09:25:22 +00:00
|
|
|
url = f"{self.ES_URL}/ta_{index_name}"
|
2021-10-28 08:49:58 +00:00
|
|
|
response = requests.get(url, auth=self.ES_AUTH)
|
2021-09-05 17:10:14 +00:00
|
|
|
exists = response.ok
|
|
|
|
|
|
|
|
if exists:
|
2021-09-21 09:25:22 +00:00
|
|
|
details = response.json()[f"ta_{index_name}"]
|
2021-09-05 17:10:14 +00:00
|
|
|
else:
|
|
|
|
details = False
|
|
|
|
|
|
|
|
return exists, details
|
|
|
|
|
|
|
|
def validate(self):
|
|
|
|
"""
|
|
|
|
check if all expected mappings and settings match
|
|
|
|
returns True when rebuild is needed
|
|
|
|
"""
|
|
|
|
|
|
|
|
if self.expected_map:
|
|
|
|
rebuild = self.validate_mappings()
|
|
|
|
if rebuild:
|
|
|
|
return rebuild
|
|
|
|
|
|
|
|
if self.expected_set:
|
|
|
|
rebuild = self.validate_settings()
|
|
|
|
if rebuild:
|
|
|
|
return rebuild
|
|
|
|
|
|
|
|
return False
|
|
|
|
|
|
|
|
def validate_mappings(self):
|
2021-09-21 09:25:22 +00:00
|
|
|
"""check if all mappings are as expected"""
|
2021-09-05 17:10:14 +00:00
|
|
|
|
|
|
|
expected_map = self.expected_map
|
2021-09-21 09:25:22 +00:00
|
|
|
now_map = self.details["mappings"]["properties"]
|
2021-09-05 17:10:14 +00:00
|
|
|
|
|
|
|
for key, value in expected_map.items():
|
|
|
|
# nested
|
2021-09-21 09:25:22 +00:00
|
|
|
if list(value.keys()) == ["properties"]:
|
|
|
|
for key_n, value_n in value["properties"].items():
|
|
|
|
if key_n not in now_map[key]["properties"].keys():
|
2021-09-05 17:10:14 +00:00
|
|
|
print(key_n, value_n)
|
|
|
|
return True
|
2021-09-21 09:25:22 +00:00
|
|
|
if not value_n == now_map[key]["properties"][key_n]:
|
2021-09-05 17:10:14 +00:00
|
|
|
print(key_n, value_n)
|
|
|
|
return True
|
|
|
|
|
|
|
|
continue
|
|
|
|
|
|
|
|
# not nested
|
|
|
|
if key not in now_map.keys():
|
|
|
|
print(key, value)
|
|
|
|
return True
|
|
|
|
if not value == now_map[key]:
|
|
|
|
print(key, value)
|
|
|
|
return True
|
|
|
|
|
|
|
|
return False
|
|
|
|
|
|
|
|
def validate_settings(self):
|
2021-09-21 09:25:22 +00:00
|
|
|
"""check if all settings are as expected"""
|
2021-09-05 17:10:14 +00:00
|
|
|
|
2021-09-21 09:25:22 +00:00
|
|
|
now_set = self.details["settings"]["index"]
|
2021-09-05 17:10:14 +00:00
|
|
|
|
|
|
|
for key, value in self.expected_set.items():
|
|
|
|
if key not in now_set.keys():
|
|
|
|
print(key, value)
|
|
|
|
return True
|
|
|
|
|
|
|
|
if not value == now_set[key]:
|
|
|
|
print(key, value)
|
|
|
|
return True
|
|
|
|
|
|
|
|
return False
|
|
|
|
|
|
|
|
def rebuild_index(self):
|
2021-09-21 09:25:22 +00:00
|
|
|
"""rebuild with new mapping"""
|
2021-09-05 17:10:14 +00:00
|
|
|
# backup
|
2021-09-21 09:25:22 +00:00
|
|
|
self.reindex("backup")
|
2021-09-05 17:10:14 +00:00
|
|
|
# delete original
|
|
|
|
self.delete_index(backup=False)
|
|
|
|
# create new
|
|
|
|
self.create_blank()
|
2021-09-21 09:25:22 +00:00
|
|
|
self.reindex("restore")
|
2021-09-05 17:10:14 +00:00
|
|
|
# delete backup
|
|
|
|
self.delete_index()
|
|
|
|
|
|
|
|
def reindex(self, method):
|
2021-09-21 09:25:22 +00:00
|
|
|
"""create on elastic search"""
|
2021-09-05 17:10:14 +00:00
|
|
|
index_name = self.index_name
|
2021-09-21 09:25:22 +00:00
|
|
|
if method == "backup":
|
|
|
|
source = f"ta_{index_name}"
|
|
|
|
destination = f"ta_{index_name}_backup"
|
|
|
|
elif method == "restore":
|
|
|
|
source = f"ta_{index_name}_backup"
|
|
|
|
destination = f"ta_{index_name}"
|
|
|
|
|
|
|
|
query = {"source": {"index": source}, "dest": {"index": destination}}
|
2021-09-05 17:10:14 +00:00
|
|
|
data = json.dumps(query)
|
2021-09-21 09:25:22 +00:00
|
|
|
url = self.ES_URL + "/_reindex?refresh=true"
|
2021-10-28 08:49:58 +00:00
|
|
|
response = requests.post(
|
|
|
|
url=url, data=data, headers=self.HEADERS, auth=self.ES_AUTH
|
|
|
|
)
|
2021-09-05 17:10:14 +00:00
|
|
|
if not response.ok:
|
|
|
|
print(response.text)
|
|
|
|
|
|
|
|
def delete_index(self, backup=True):
|
2021-09-21 09:25:22 +00:00
|
|
|
"""delete index passed as argument"""
|
2021-09-05 17:10:14 +00:00
|
|
|
if backup:
|
2021-09-21 09:25:22 +00:00
|
|
|
url = f"{self.ES_URL}/ta_{self.index_name}_backup"
|
2021-09-05 17:10:14 +00:00
|
|
|
else:
|
2021-09-21 09:25:22 +00:00
|
|
|
url = f"{self.ES_URL}/ta_{self.index_name}"
|
2021-10-28 08:49:58 +00:00
|
|
|
response = requests.delete(url, auth=self.ES_AUTH)
|
2021-09-05 17:10:14 +00:00
|
|
|
if not response.ok:
|
|
|
|
print(response.text)
|
|
|
|
|
|
|
|
def create_blank(self):
|
2021-09-21 09:25:22 +00:00
|
|
|
"""apply new mapping and settings for blank new index"""
|
2021-09-05 17:10:14 +00:00
|
|
|
expected_map = self.expected_map
|
|
|
|
expected_set = self.expected_set
|
|
|
|
# stich payload
|
|
|
|
payload = {}
|
|
|
|
if expected_set:
|
|
|
|
payload.update({"settings": expected_set})
|
|
|
|
if expected_map:
|
|
|
|
payload.update({"mappings": {"properties": expected_map}})
|
|
|
|
# create
|
2021-09-21 09:25:22 +00:00
|
|
|
url = f"{self.ES_URL}/ta_{self.index_name}"
|
2021-09-05 17:10:14 +00:00
|
|
|
data = json.dumps(payload)
|
2021-10-28 08:49:58 +00:00
|
|
|
response = requests.put(
|
|
|
|
url=url, data=data, headers=self.HEADERS, auth=self.ES_AUTH
|
|
|
|
)
|
2021-09-05 17:10:14 +00:00
|
|
|
if not response.ok:
|
|
|
|
print(response.text)
|
|
|
|
|
|
|
|
|
|
|
|
class ElasticBackup:
|
2021-09-21 09:25:22 +00:00
|
|
|
"""dump index to nd-json files for later bulk import"""
|
2021-09-05 17:10:14 +00:00
|
|
|
|
|
|
|
def __init__(self, index_config):
|
|
|
|
self.config = AppConfig().config
|
|
|
|
self.index_config = index_config
|
2021-09-21 09:25:22 +00:00
|
|
|
self.timestamp = datetime.now().strftime("%Y%m%d")
|
2021-09-16 10:34:20 +00:00
|
|
|
self.backup_files = []
|
2021-09-05 17:10:14 +00:00
|
|
|
|
|
|
|
def get_all_documents(self, index_name):
|
2021-09-21 09:25:22 +00:00
|
|
|
"""export all documents of a single index"""
|
|
|
|
headers = {"Content-type": "application/json"}
|
|
|
|
es_url = self.config["application"]["es_url"]
|
2021-10-28 08:49:58 +00:00
|
|
|
es_auth = self.config["application"]["es_auth"]
|
2021-09-05 17:10:14 +00:00
|
|
|
# get PIT ID
|
2021-09-21 09:25:22 +00:00
|
|
|
url = f"{es_url}/ta_{index_name}/_pit?keep_alive=1m"
|
2021-10-28 08:49:58 +00:00
|
|
|
response = requests.post(url, auth=es_auth)
|
2021-09-05 17:10:14 +00:00
|
|
|
json_data = json.loads(response.text)
|
2021-09-21 09:25:22 +00:00
|
|
|
pit_id = json_data["id"]
|
2021-09-05 17:10:14 +00:00
|
|
|
# build query
|
|
|
|
data = {
|
|
|
|
"query": {"match_all": {}},
|
2021-09-21 09:25:22 +00:00
|
|
|
"size": 100,
|
|
|
|
"pit": {"id": pit_id, "keep_alive": "1m"},
|
|
|
|
"sort": [{"_id": {"order": "asc"}}],
|
2021-09-05 17:10:14 +00:00
|
|
|
}
|
|
|
|
query_str = json.dumps(data)
|
2021-09-21 09:25:22 +00:00
|
|
|
url = es_url + "/_search"
|
2021-09-05 17:10:14 +00:00
|
|
|
# loop until nothing left
|
|
|
|
all_results = []
|
|
|
|
while True:
|
2021-10-28 08:49:58 +00:00
|
|
|
response = requests.get(
|
|
|
|
url, data=query_str, headers=headers, auth=es_auth
|
|
|
|
)
|
2021-09-05 17:10:14 +00:00
|
|
|
json_data = json.loads(response.text)
|
2021-09-21 09:25:22 +00:00
|
|
|
all_hits = json_data["hits"]["hits"]
|
2021-09-05 17:10:14 +00:00
|
|
|
if all_hits:
|
|
|
|
for hit in all_hits:
|
2021-09-21 09:25:22 +00:00
|
|
|
search_after = hit["sort"]
|
2021-09-05 17:10:14 +00:00
|
|
|
all_results.append(hit)
|
|
|
|
# update search_after with last hit data
|
2021-09-21 09:25:22 +00:00
|
|
|
data["search_after"] = search_after
|
2021-09-05 17:10:14 +00:00
|
|
|
query_str = json.dumps(data)
|
|
|
|
else:
|
|
|
|
break
|
|
|
|
# clean up PIT
|
|
|
|
query_str = json.dumps({"id": pit_id})
|
2021-10-28 08:49:58 +00:00
|
|
|
requests.delete(
|
|
|
|
es_url + "/_pit", data=query_str, headers=headers, auth=es_auth
|
|
|
|
)
|
2021-09-05 17:10:14 +00:00
|
|
|
|
|
|
|
return all_results
|
|
|
|
|
|
|
|
@staticmethod
|
|
|
|
def build_bulk(all_results):
|
2021-09-21 09:25:22 +00:00
|
|
|
"""build bulk query data from all_results"""
|
2021-09-05 17:10:14 +00:00
|
|
|
bulk_list = []
|
|
|
|
|
|
|
|
for document in all_results:
|
2021-09-21 09:25:22 +00:00
|
|
|
document_id = document["_id"]
|
|
|
|
es_index = document["_index"]
|
2021-09-16 10:34:20 +00:00
|
|
|
action = {"index": {"_index": es_index, "_id": document_id}}
|
2021-09-21 09:25:22 +00:00
|
|
|
source = document["_source"]
|
2021-09-05 17:10:14 +00:00
|
|
|
bulk_list.append(json.dumps(action))
|
|
|
|
bulk_list.append(json.dumps(source))
|
|
|
|
|
|
|
|
# add last newline
|
2021-09-21 09:25:22 +00:00
|
|
|
bulk_list.append("\n")
|
|
|
|
file_content = "\n".join(bulk_list)
|
2021-09-05 17:10:14 +00:00
|
|
|
|
|
|
|
return file_content
|
|
|
|
|
2021-09-16 10:34:20 +00:00
|
|
|
def write_es_json(self, file_content, index_name):
|
2021-09-21 09:25:22 +00:00
|
|
|
"""write nd-json file for es _bulk API to disk"""
|
|
|
|
cache_dir = self.config["application"]["cache_dir"]
|
|
|
|
file_name = f"es_{index_name}-{self.timestamp}.json"
|
|
|
|
file_path = os.path.join(cache_dir, "backup", file_name)
|
|
|
|
with open(file_path, "w", encoding="utf-8") as f:
|
2021-09-16 10:34:20 +00:00
|
|
|
f.write(file_content)
|
|
|
|
|
|
|
|
self.backup_files.append(file_path)
|
|
|
|
|
|
|
|
def write_ta_json(self, all_results, index_name):
|
2021-09-21 09:25:22 +00:00
|
|
|
"""write generic json file to disk"""
|
|
|
|
cache_dir = self.config["application"]["cache_dir"]
|
|
|
|
file_name = f"ta_{index_name}-{self.timestamp}.json"
|
|
|
|
file_path = os.path.join(cache_dir, "backup", file_name)
|
|
|
|
to_write = [i["_source"] for i in all_results]
|
2021-09-16 10:34:20 +00:00
|
|
|
file_content = json.dumps(to_write)
|
2021-09-21 09:25:22 +00:00
|
|
|
with open(file_path, "w", encoding="utf-8") as f:
|
2021-09-05 17:10:14 +00:00
|
|
|
f.write(file_content)
|
|
|
|
|
2021-09-16 10:34:20 +00:00
|
|
|
self.backup_files.append(file_path)
|
|
|
|
|
|
|
|
def zip_it(self):
|
2021-09-21 09:25:22 +00:00
|
|
|
"""pack it up into single zip file"""
|
|
|
|
cache_dir = self.config["application"]["cache_dir"]
|
|
|
|
file_name = f"ta_backup-{self.timestamp}.zip"
|
|
|
|
backup_folder = os.path.join(cache_dir, "backup")
|
2021-09-18 10:10:16 +00:00
|
|
|
backup_file = os.path.join(backup_folder, file_name)
|
2021-09-16 10:34:20 +00:00
|
|
|
|
|
|
|
with zipfile.ZipFile(
|
2021-09-21 09:25:22 +00:00
|
|
|
backup_file, "w", compression=zipfile.ZIP_DEFLATED
|
|
|
|
) as zip_f:
|
2021-09-16 10:34:20 +00:00
|
|
|
for backup_file in self.backup_files:
|
2021-09-18 10:10:16 +00:00
|
|
|
zip_f.write(backup_file, os.path.basename(backup_file))
|
2021-09-16 10:34:20 +00:00
|
|
|
|
|
|
|
# cleanup
|
|
|
|
for backup_file in self.backup_files:
|
|
|
|
os.remove(backup_file)
|
|
|
|
|
2021-09-05 17:10:14 +00:00
|
|
|
def post_bulk_restore(self, file_name):
|
2021-09-21 09:25:22 +00:00
|
|
|
"""send bulk to es"""
|
|
|
|
cache_dir = self.config["application"]["cache_dir"]
|
|
|
|
es_url = self.config["application"]["es_url"]
|
2021-10-28 08:49:58 +00:00
|
|
|
es_auth = self.config["application"]["es_auth"]
|
2021-09-21 09:25:22 +00:00
|
|
|
headers = {"Content-type": "application/x-ndjson"}
|
2021-09-05 17:10:14 +00:00
|
|
|
file_path = os.path.join(cache_dir, file_name)
|
|
|
|
|
2021-09-21 09:25:22 +00:00
|
|
|
with open(file_path, "r", encoding="utf-8") as f:
|
2021-09-05 17:10:14 +00:00
|
|
|
query_str = f.read()
|
|
|
|
|
2021-09-20 12:10:39 +00:00
|
|
|
if not query_str.strip():
|
|
|
|
return
|
|
|
|
|
2021-09-21 09:25:22 +00:00
|
|
|
url = es_url + "/_bulk"
|
2021-10-28 08:49:58 +00:00
|
|
|
request = requests.post(
|
|
|
|
url, data=query_str, headers=headers, auth=es_auth
|
|
|
|
)
|
2021-09-05 17:10:14 +00:00
|
|
|
if not request.ok:
|
|
|
|
print(request.text)
|
|
|
|
|
2021-09-18 10:10:16 +00:00
|
|
|
def unpack_zip_backup(self):
|
2021-09-21 09:25:22 +00:00
|
|
|
"""extract backup zip and return filelist"""
|
|
|
|
cache_dir = self.config["application"]["cache_dir"]
|
|
|
|
backup_dir = os.path.join(cache_dir, "backup")
|
2021-09-25 11:59:54 +00:00
|
|
|
backup_files = os.listdir(backup_dir)
|
|
|
|
all_backup_files = ignore_filelist(backup_files)
|
2021-09-05 17:10:14 +00:00
|
|
|
all_available_backups = [
|
2021-09-21 09:25:22 +00:00
|
|
|
i
|
2021-09-25 11:59:54 +00:00
|
|
|
for i in all_backup_files
|
2021-09-21 09:25:22 +00:00
|
|
|
if i.startswith("ta_") and i.endswith(".zip")
|
2021-09-05 17:10:14 +00:00
|
|
|
]
|
2021-09-18 10:10:16 +00:00
|
|
|
all_available_backups.sort()
|
|
|
|
newest_backup = all_available_backups[-1]
|
|
|
|
file_path = os.path.join(backup_dir, newest_backup)
|
|
|
|
|
2021-09-21 09:25:22 +00:00
|
|
|
with zipfile.ZipFile(file_path, "r") as z:
|
2021-09-18 10:10:16 +00:00
|
|
|
zip_content = z.namelist()
|
|
|
|
z.extractall(backup_dir)
|
|
|
|
|
|
|
|
return zip_content
|
|
|
|
|
|
|
|
def restore_json_files(self, zip_content):
|
2021-09-21 09:25:22 +00:00
|
|
|
"""go through the unpacked files and restore"""
|
2021-09-18 10:10:16 +00:00
|
|
|
|
2021-09-21 09:25:22 +00:00
|
|
|
cache_dir = self.config["application"]["cache_dir"]
|
|
|
|
backup_dir = os.path.join(cache_dir, "backup")
|
2021-09-18 10:10:16 +00:00
|
|
|
|
|
|
|
for json_f in zip_content:
|
|
|
|
|
|
|
|
file_name = os.path.join(backup_dir, json_f)
|
|
|
|
|
2021-09-21 09:25:22 +00:00
|
|
|
if not json_f.startswith("es_") or not json_f.endswith(".json"):
|
2021-09-18 10:10:16 +00:00
|
|
|
os.remove(file_name)
|
|
|
|
continue
|
|
|
|
|
2021-09-21 09:25:22 +00:00
|
|
|
print("restoring: " + json_f)
|
2021-09-05 17:10:14 +00:00
|
|
|
self.post_bulk_restore(file_name)
|
2021-09-18 10:10:16 +00:00
|
|
|
os.remove(file_name)
|
2021-09-05 17:10:14 +00:00
|
|
|
|
2021-11-26 09:19:18 +00:00
|
|
|
def index_exists(self, index_name):
|
|
|
|
"""check if index already exists to skip"""
|
|
|
|
es_url = self.config["application"]["es_url"]
|
|
|
|
es_auth = self.config["application"]["es_auth"]
|
|
|
|
url = f"{es_url}/ta_{index_name}"
|
|
|
|
response = requests.get(url, auth=es_auth)
|
|
|
|
|
|
|
|
return response.ok
|
|
|
|
|
2021-09-05 17:10:14 +00:00
|
|
|
|
|
|
|
def backup_all_indexes():
|
2021-09-21 09:25:22 +00:00
|
|
|
"""backup all es indexes to disk"""
|
2021-09-05 17:10:14 +00:00
|
|
|
backup_handler = ElasticBackup(INDEX_CONFIG)
|
|
|
|
|
|
|
|
for index in backup_handler.index_config:
|
2021-09-21 09:25:22 +00:00
|
|
|
index_name = index["index_name"]
|
2021-11-26 09:19:18 +00:00
|
|
|
if not backup_handler.index_exists(index_name):
|
|
|
|
continue
|
2021-09-05 17:10:14 +00:00
|
|
|
all_results = backup_handler.get_all_documents(index_name)
|
|
|
|
file_content = backup_handler.build_bulk(all_results)
|
2021-09-16 10:34:20 +00:00
|
|
|
backup_handler.write_es_json(file_content, index_name)
|
|
|
|
backup_handler.write_ta_json(all_results, index_name)
|
|
|
|
|
|
|
|
backup_handler.zip_it()
|
2021-09-05 17:10:14 +00:00
|
|
|
|
|
|
|
|
|
|
|
def restore_from_backup():
|
2021-09-21 09:25:22 +00:00
|
|
|
"""restore indexes from backup file"""
|
2021-09-05 17:10:14 +00:00
|
|
|
# delete
|
|
|
|
index_check(force_restore=True)
|
|
|
|
# recreate
|
|
|
|
backup_handler = ElasticBackup(INDEX_CONFIG)
|
2021-09-18 10:10:16 +00:00
|
|
|
zip_content = backup_handler.unpack_zip_backup()
|
|
|
|
backup_handler.restore_json_files(zip_content)
|
2021-09-05 17:10:14 +00:00
|
|
|
|
|
|
|
|
|
|
|
def index_check(force_restore=False):
|
2021-09-21 09:25:22 +00:00
|
|
|
"""check if all indexes are created and have correct mapping"""
|
2021-09-08 05:31:33 +00:00
|
|
|
|
|
|
|
backed_up = False
|
|
|
|
|
2021-09-05 17:10:14 +00:00
|
|
|
for index in INDEX_CONFIG:
|
2021-09-21 09:25:22 +00:00
|
|
|
index_name = index["index_name"]
|
|
|
|
expected_map = index["expected_map"]
|
|
|
|
expected_set = index["expected_set"]
|
2021-09-05 17:10:14 +00:00
|
|
|
handler = ElasticIndex(index_name, expected_map, expected_set)
|
|
|
|
# force restore
|
|
|
|
if force_restore:
|
|
|
|
handler.delete_index(backup=False)
|
|
|
|
handler.create_blank()
|
|
|
|
continue
|
|
|
|
|
|
|
|
# create new
|
|
|
|
if not handler.exists:
|
2021-09-21 09:25:22 +00:00
|
|
|
print(f"create new blank index with name ta_{index_name}...")
|
2021-09-05 17:10:14 +00:00
|
|
|
handler.create_blank()
|
|
|
|
continue
|
|
|
|
|
|
|
|
# validate index
|
|
|
|
rebuild = handler.validate()
|
|
|
|
if rebuild:
|
2021-09-08 05:31:33 +00:00
|
|
|
# make backup before rebuild
|
|
|
|
if not backed_up:
|
2021-09-21 09:25:22 +00:00
|
|
|
print("running backup first")
|
2021-09-08 05:31:33 +00:00
|
|
|
backup_all_indexes()
|
|
|
|
backed_up = True
|
|
|
|
|
2021-09-21 09:25:22 +00:00
|
|
|
print(f"applying new mappings to index ta_{index_name}...")
|
2021-09-05 17:10:14 +00:00
|
|
|
handler.rebuild_index()
|
|
|
|
continue
|
|
|
|
|
|
|
|
# else all good
|
2021-09-21 09:25:22 +00:00
|
|
|
print(f"ta_{index_name} index is created and up to date...")
|