Repository navigation
Expand file tree
/
Copy pathload_dataset.py
More file actions
128 lines (110 loc) · 3.76 KB
/
Copy pathload_dataset.py
File metadata and controls
128 lines (110 loc) · 3.76 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
import argparse
import json
import logging
import sys
from pathlib import Path
from elasticsearch import Elasticsearch, ApiError
from elasticsearch.helpers import streaming_bulk
logger = logging.getLogger(__name__)
def get_file_basename(file: str):
"""Return the basename of a file path."""
return Path(file).stem
def load_dataset(file: str, index_name: str):
"""Yield non-empty NDJSON lines as Elasticsearch bulk indexing actions."""
with open(file, "r") as f:
for line in f:
if line.strip():
yield {
"_index": index_name,
"_source": json.loads(line),
}
def bulk_upload(index_name: str, file: str, chunk_size: int, index_settings: dict | None):
"""Create an index if needed and upload dataset documents in bulk."""
client = Elasticsearch("http://localhost:9200")
if not client.indices.exists(index=index_name):
try:
client.indices.create(index=index_name, body=index_settings)
except ApiError as error:
logger.error("Index creation failed: %s", error)
sys.exit(1)
logger.info("Index '%s' created", index_name)
else:
logger.warning("Index '%s' already exists, skipping index creation", index_name)
logger.info("Indexing documents from '%s' into '%s' index", file, index_name)
processed_docs: int = 0
indexed_docs: int = 0
failed_docs: int = 0
for success, error in streaming_bulk(
client=client,
actions=load_dataset(file=file, index_name=index_name),
chunk_size=chunk_size,
raise_on_error=False
):
processed_docs += 1
if not success:
failed_docs += 1
logger.error("Failed to load document into index '%s': %s", index_name, error)
else:
indexed_docs += 1
logger.info(
"Indexing complete: %d processed, %d indexed, %d failed",
processed_docs,
indexed_docs,
failed_docs,
)
def main():
parser = argparse.ArgumentParser(description="Load dataset data into an ElasticSearch index")
parser.add_argument(
"file",
type=str,
help="Path to dataset file"
)
parser.add_argument(
"-i", "--index-name",
dest="index_name",
default=None,
type=str,
help="Name of the index"
)
parser.add_argument(
"-s", "--settings-file",
dest="settings_file",
type=str,
help="Path to index settings file (includes settings and/or mappings)",
)
parser.add_argument(
"-c", "--chunk-size",
dest="chunk_size",
default=1000,
type=int,
help="Number of documents to send per bulk batch (default: 1000)"
)
args = parser.parse_args()
logging.getLogger("elastic_transport").setLevel(logging.WARNING)
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s %(levelname)s: %(message)s",
datefmt="%Y-%m-%dT%H:%M:%S",
)
settings_file: str = args.settings_file
settings: dict | None = None
if settings_file:
try:
with open(settings_file, "r") as f:
settings = json.load(f)
except FileNotFoundError:
logger.error("Index settings file not found, aborting")
sys.exit(1)
except json.JSONDecodeError as error:
logger.error("Index settings file parsing failed: %s", error)
sys.exit(1)
if args.chunk_size <= 0:
raise ValueError("--chunk-size flag must be a positive integer")
bulk_upload(
file=args.file,
index_name=args.index_name or get_file_basename(args.file),
index_settings=settings,
chunk_size=args.chunk_size
)
if __name__ == "__main__":
main()