Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions config.ini
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,9 @@ omid=oc_index.identifier.omid:OMIDManager
verbose=1
logdir=logs

[dump]
output=_out_

# Redis data source info
[redis]
host=127.0.0.1
Expand Down
2 changes: 1 addition & 1 deletion oc_index/scripts/cits2redis.py
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@
from oc_index.utils.logging import get_logger
from oc_index.utils.config import get_config

csv.field_size_limit(sys.maxsize)
csv.field_size_limit(sys.maxsize) if os.name != "nt" else csv.field_size_limit(2**31 - 1)

NEEDLE = "ci/"
BATCH_SIZE = 50_000
Expand Down
9 changes: 6 additions & 3 deletions oc_index/scripts/cnc.py
Original file line number Diff line number Diff line change
Expand Up @@ -400,15 +400,15 @@ def main():

redis_br = redis.Redis(
host="127.0.0.1",
port=6379,
port=int(_config.get("redis", "port")),
db=int(_config.get("cnc", "db_br")),
decode_responses=True,
)
redis_cits_cache = redis.Redis(host="127.0.0.1", port=6379, db=int(_config.get("cnc", "db_omid")))
redis_cits_cache = redis.Redis(host="127.0.0.1", port=int(_config.get("redis", "port")), db=int(_config.get("cnc", "db_omid")))

redis_cits = redis.Redis(
host="127.0.0.1",
port=6379,
port=int(_config.get("redis", "port")),
db=int(_config.get("cnc", "db_cits")),
decode_responses=True,
)
Expand Down Expand Up @@ -454,3 +454,6 @@ def main():
# 4. Continue with the rest of your code **after all files are done**
# e.g., merging outputs, generating RDF/CSV summary, logging, etc.
# >> post_processing(output_dir)

if __name__ == "__main__":
main()
21 changes: 13 additions & 8 deletions oc_index/scripts/dump_index.py
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@
source: str
service_name: str
index_identifier: str
file_output_dir: str


def zip_and_cleanup(csv_dir, rdf_dir, slx_dir, files_per_zip, force = False, pnum=1):
Expand Down Expand Up @@ -99,7 +100,7 @@ def chunk_list(lst, n):

def main():

global _logger, idbase_url, baseurl, agent, source, service_name, index_identifier
global _logger, idbase_url, baseurl, agent, source, service_name, index_identifier, file_output_dir
global FILE_OUTPUT_DIR

arg_parser = ArgumentParser(description="Dump OpenCitations Index data. This process reads all the data in Redis and creates a new data dump for the OpenCitations Index. The outputs are compressed, to all dump formats: CSV, RDF, SCHOLIX. **Make sure the Redis datasets are populated before running this script**")
Expand Down Expand Up @@ -132,6 +133,7 @@ def main():
source = _config.get("INDEX", "source")
service_name = _config.get("INDEX", "service")
index_identifier = _config.get("INDEX", "identifier")
file_output_dir = _config.get("dump", "output")

dump_date = datetime.now().strftime("%Y%m%d")
if args.date:
Expand Down Expand Up @@ -165,12 +167,13 @@ def main():
# === REDIS ===
REDIS_CITS_DB = _config.get("cnc", "db_cits")
REDIS_METADATA_DB = _config.get("INDEX", "db")
REDIS_PORT = int(_config.get("redis", "port"))

redis_cits = redis.Redis(host='localhost', port=6379, db=int(REDIS_CITS_DB), decode_responses=True)
redis_cits = redis.Redis(host='localhost', port=REDIS_PORT, db=int(REDIS_CITS_DB), decode_responses=True)
# Sample data of redis_cits:
# "06304836421": "[\"06290442260\", \"0606973973\", \"06290442260\", \"061204315925\"]"

redis_metadata = redis.Redis(host='localhost', port=6379, db=int(REDIS_METADATA_DB), decode_responses=True)
redis_metadata = redis.Redis(host='localhost', port=REDIS_PORT, db=int(REDIS_METADATA_DB), decode_responses=True)
# Sample data of redis_metadata:
# "omid:br/061601556475": "{\"date\": \"2019\", \"valid\": true, \"orcid\": [\"0000-0002-6819-0387\"], \"issn\": [\"0886-022X\", \"1525-6049\"]}"

Expand Down Expand Up @@ -231,7 +234,7 @@ def main():
processes = []
for idx,chunk in enumerate(g_chunks):
_logger.info(f" - Process {idx} elaborates #citations in chunk: {len(chunk)}")
p = Process(target=process_pair, args=(chunk, idx, br_meta, cursor == 0))
p = Process(target=process_pair, args=(chunk, idx, br_meta, cursor == 0, baseurl, file_output_dir))
p.start()
processes.append(p)

Expand Down Expand Up @@ -268,8 +271,7 @@ def main():
break


def process_pair(pairs, pnum, br_meta, end_cursor = False):

def process_pair(pairs, pnum, br_meta, end_cursor = False, baseurl: str = "", file_output_dir: str = ""):
data_to_dump = []

for pair in pairs:
Expand Down Expand Up @@ -321,10 +323,10 @@ def process_pair(pairs, pnum, br_meta, end_cursor = False):

# write p_data_to_dump to files when range CITATIONS_PER_FILE is reached
# if len(data_to_dump[pnum]) >= CITATIONS_PER_FILE or end_cursor:
_logger.info(f"Storing {len(data_to_dump)} citations data of task {pnum}...")
#_logger.info(f"Storing {len(data_to_dump)} citations data of task {pnum}...")
# write to files
index_ts_storer = CitationStorer(
FILE_OUTPUT_DIR,
file_output_dir,
baseurl + "/" if not baseurl.endswith("/") else baseurl,
store_as=["csv_data","rdf_data","scholix_data"],
suffix= str(pnum)
Expand All @@ -343,3 +345,6 @@ def process_pair(pairs, pnum, br_meta, end_cursor = False):
force = end_cursor,
pnum = pnum
)

if __name__ == "__main__":
main()
4 changes: 2 additions & 2 deletions oc_index/scripts/meta2redis.py
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,7 @@
from oc_index.utils.config import get_config

console = Console()
csv.field_size_limit(sys.maxsize)
csv.field_size_limit(sys.maxsize) if os.name != "nt" else csv.field_size_limit(2**31 - 1)

BASE_IRI = "https://w3id.org/oc/meta/"
DIR_SPLIT = 10000
Expand Down Expand Up @@ -316,7 +316,7 @@ def _process_csv_file(a_csv_file, rconn_db_br, rconn_db_ra, rconn_db_metadata):
ra_data = defaultdict(set)
metadata = {}

text_file = io.TextIOWrapper(a_csv_file)
text_file = io.TextIOWrapper(a_csv_file, encoding="utf-8")
try:
for o_row in csv.DictReader(text_file):
br_ids = get_key_ids(o_row["id"])
Expand Down