From 4a6c8cd70f546cad9c3a85a738d02fe07f8b4c5b Mon Sep 17 00:00:00 2001 From: krzywonos <98768745+krzywonos@users.noreply.github.com> Date: Mon, 20 Oct 2025 12:40:18 +0200 Subject: [PATCH 1/4] meta2redis.py - Windows compatibility and fixes for reading CSV files - field size limit for CSV files set to C 32-bit int limit when running on Windows - specified conversion from binary to utf-8 for TextIOWrapper instead of using default system settings - CSV files are now read as binary to conform to the usage of TextIOWrapper - added main function call --- scripts/meta2redis.py | 15 +++++++++++---- 1 file changed, 11 insertions(+), 4 deletions(-) diff --git a/scripts/meta2redis.py b/scripts/meta2redis.py index f49b5d2..b1c525f 100755 --- a/scripts/meta2redis.py +++ b/scripts/meta2redis.py @@ -31,7 +31,11 @@ from oc.index.utils.config import get_config _config = get_config() -csv.field_size_limit(sys.maxsize) + +if os.name != "nt": + csv.field_size_limit(sys.maxsize) +else: + csv.field_size_limit(2**31 - 1) # glob indexes br_ids = _config.get("cnc", "br_ids").split(",") @@ -97,7 +101,7 @@ def _p_csvfile(a_csv_file,csv_name,rconn_db_br, rconn_db_ra, rconn_db_metadata): db_ra_buffer = [] db_metadata_buffer = [] - l_brs = list(csv.DictReader(io.TextIOWrapper(a_csv_file))) + l_brs = list(csv.DictReader(io.TextIOWrapper(a_csv_file, encoding="utf-8"))) # walk through each citation in the CSV logger.info("Walking through all the "+str( len(l_brs) )+" BRs (rows) in: "+str(csv_name) ) @@ -203,7 +207,7 @@ def upload2redis(dump_path="", redishost="localhost", redisport="6379", redisbat # Handle single CSV file csv_name = os.path.basename(dump_path) logger.info(f"CSV: Processing direct CSV file: {csv_name}") - with open(dump_path, 'r', encoding='utf-8') as csv_file: + with open(dump_path, 'rb') as csv_file: _p_csvfile(csv_file,csv_name, rconn_db_br, rconn_db_ra, rconn_db_metadata) else: logger.warning(f"Unsupported file type: {dump_path}") @@ -219,7 +223,7 @@ def upload2redis(dump_path="", redishost="localhost", redisport="6379", redisbat if filename.endswith(".csv"): logger.info(f"CSV: Processing direct CSV file: {filename}") - with open(filepath, 'r', encoding='utf-8') as csv_file: + with open(filepath, 'rb') as csv_file: _p_csvfile(csv_file, filename, rconn_db_br, rconn_db_ra, rconn_db_metadata) else: logger.error(f"Path does not exist or is neither a file nor directory: {dump_path}") @@ -265,3 +269,6 @@ def main(): ) logger.info("A total of unique "+str(res[0])+" BR OMIDs and "+str(res[1])+" RA OMIDs have been found and added to Redis.") + +if __name__ == "__main__": + main() From 66d930909f544707847c2d5078238b49784a959c Mon Sep 17 00:00:00 2001 From: Hubert Krzywonos Date: Wed, 9 Sep 2026 11:32:51 +0200 Subject: [PATCH 2/4] meta2redis/cits2redis - minor fix for Windows CSV limits --- oc_index/scripts/cits2redis.py | 2 +- oc_index/scripts/meta2redis.py | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/oc_index/scripts/cits2redis.py b/oc_index/scripts/cits2redis.py index e9b5834..019d61d 100755 --- a/oc_index/scripts/cits2redis.py +++ b/oc_index/scripts/cits2redis.py @@ -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 diff --git a/oc_index/scripts/meta2redis.py b/oc_index/scripts/meta2redis.py index 6dda291..60ed4c8 100755 --- a/oc_index/scripts/meta2redis.py +++ b/oc_index/scripts/meta2redis.py @@ -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 From 298d12a079c17d9e5b468bb9cb5e394752bf2d24 Mon Sep 17 00:00:00 2001 From: Hubert Krzywonos Date: Wed, 9 Sep 2026 11:39:42 +0200 Subject: [PATCH 3/4] deleted extra meta2redis file leftover from a merging mistake --- scripts/meta2redis.py | 274 ------------------------------------------ 1 file changed, 274 deletions(-) delete mode 100755 scripts/meta2redis.py diff --git a/scripts/meta2redis.py b/scripts/meta2redis.py deleted file mode 100755 index b1c525f..0000000 --- a/scripts/meta2redis.py +++ /dev/null @@ -1,274 +0,0 @@ -#!python -# Copyright (c) 2023 Ivan Heibi. -# -# Permission to use, copy, modify, and/or distribute this software for any purpose -# with or without fee is hereby granted, provided that the above copyright notice -# and this permission notice appear in all copies. -# -# THE SOFTWARE IS PROVIDED "AS IS" AND THE AUTHOR DISCLAIMS ALL WARRANTIES WITH -# REGARD TO THIS SOFTWARE INCLUDING ALL IMPLIED WARRANTIES OF MERCHANTABILITY AND -# FITNESS. IN NO EVENT SHALL THE AUTHOR BE LIABLE FOR ANY SPECIAL, DIRECT, INDIRECT, -# OR CONSEQUENTIAL DAMAGES OR ANY DAMAGES WHATSOEVER RESULTING FROM LOSS OF USE, -# DATA OR PROFITS, WHETHER IN AN ACTION OF CONTRACT, NEGLIGENCE OR OTHER TORTIOUS -# ACTION, ARISING OUT OF OR IN CONNECTION WITH THE USE OR PERFORMANCE OF THIS -# SOFTWARE. - -import csv -import json -from zipfile import ZipFile -import tarfile -import os -import datetime -import io -import argparse -from redis import Redis -import re -import sys - -from tqdm import tqdm -from collections import defaultdict -from oc.index.utils.logging import get_logger -from oc.index.utils.config import get_config - -_config = get_config() - -if os.name != "nt": - csv.field_size_limit(sys.maxsize) -else: - csv.field_size_limit(2**31 - 1) - -# glob indexes -br_ids = _config.get("cnc", "br_ids").split(",") -ra_ids = _config.get("cnc", "ra_ids").split(",") - -br_index = defaultdict(set) -ra_index = defaultdict(set) - -class RedisDB(object): - - def __init__(self, redishost, redisport, redisbatchsize, _db): - self.redisbatchsize = int(redisbatchsize) - self.rconn = Redis(host=redishost, port=redisport, db=_db) - - def set_data(self, data, force=False, type=None): - if len(data) >= self.redisbatchsize or force: - for item in data: - _k = item[0] - _v = item[1] - _v_val = _v - - if type == "br": - if _k in br_index: - br_index[_k].update( set(_v) ) - else: - br_index[_k] = set(_v) - _v_val = "; ".join(br_index[_k]) - - elif type == "ra": - if _k in ra_index: - ra_index[_k].update( set(_v) ) - else: - ra_index[_k] = set(_v) - _v_val = "; ".join(ra_index[_k]) - - self.rconn.set(_k, _v_val) - - return len(data) - return 0 - -def get_key_ids(text): - return text.split(" ") - -def get_att_ids(text): - bracket_contents = re.findall(r'\[(.*?)\]', text) - return [part.split() for part in bracket_contents] - -def get_id_val(l_ids,l_id_type = []): - res = [] - for _id in l_ids: - for _id_type in l_id_type: - if _id.startswith(_id_type): - res.append(_id) - return res - -def _p_csvfile(a_csv_file,csv_name,rconn_db_br, rconn_db_ra, rconn_db_metadata): - - global _config - logger = get_logger() - - # set buffers - db_br_buffer = [] - db_ra_buffer = [] - db_metadata_buffer = [] - - l_brs = list(csv.DictReader(io.TextIOWrapper(a_csv_file, encoding="utf-8"))) - - # walk through each citation in the CSV - logger.info("Walking through all the "+str( len(l_brs) )+" BRs (rows) in: "+str(csv_name) ) - for o_row in tqdm(l_brs): - - # list of BR ids - # > update the to be added in REDIS () - br_ids = get_key_ids(o_row["id"]) - br_ids_omid = get_id_val(br_ids,["omid"]) - br_ids_other = [x for x in br_ids if x not in br_ids_omid] - # Add it to the list of BRs - for __oid in br_ids_other: - db_br_buffer.append( - ( - __oid, - br_ids_omid - ) - ) - - # list of RA ids - # > update the to be added in REDIS () - l_ra_ids = get_att_ids(o_row["author"]) # [ ["omid:123","orcid:1111-2222"], ["omid:321","orcid:3333-4444"], ...] - for _ra in l_ra_ids: # > ["omid:123","orcid:1111-2222"] - ra_ids_omid = get_id_val(_ra,["omid"]) # > ["omid:123"] - ra_ids_other = [x for x in _ra if x not in ra_ids_omid] # > ["orcid:1111-2222"] - # Add it to the list of RAs - for __oid in ra_ids_other: - db_ra_buffer.append( - ( - __oid, - ra_ids_omid - ) - ) - - - # metadata of each br - # > update the to be added in REDIS () - for _omid in br_ids_omid: - - # Take only ORCIDs - orcids = [] - for _ra in l_ra_ids: - orcids += get_id_val(_ra,["orcid"]) # > ["orcid:1111-2222"] - - # Take only ISSNs - issns = [] - l_venue_ids = get_att_ids(o_row["venue"]) - for _venue in l_venue_ids: - issns += get_id_val(_venue,["issn"]) # > ["issn:1111-2222"] - - db_metadata_buffer.append( - ( - _omid, - json.dumps( - { - "date": str(o_row["pub_date"]), - "valid": True, - "orcid": [a.replace("orcid:","") for a in orcids] if len(orcids) > 0 else [], - "issn": [a.replace("issn:","") for a in issns] if len(issns) > 0 else [] - } - ) - ) - ) - - # Set last data in Redis - logger.info("Updating Redis ... ") - rconn_db_metadata.set_data(db_metadata_buffer, True) - rconn_db_br.set_data(db_br_buffer, True, type= "br") - rconn_db_ra.set_data(db_ra_buffer, True, type= "ra") - - -def upload2redis(dump_path="", redishost="localhost", redisport="6379", redisbatchsize="10000", db_omid = "9", db_br="10", db_ra="11", db_metadata="12"): - global _config - logger = get_logger() - - #rconn_db_omid = RedisDB(redishost, redisport, redisbatchsize, db_omid) - rconn_db_br = RedisDB(redishost, redisport, redisbatchsize, db_br) - rconn_db_ra = RedisDB(redishost, redisport, redisbatchsize, db_ra) - rconn_db_metadata = RedisDB(redishost, redisport, redisbatchsize, db_metadata) - - # Check if dump_path is a single archive file - if os.path.isfile(dump_path): - if dump_path.endswith(".zip"): - # Handle single ZIP file - with ZipFile(dump_path) as archive: - logger.info(f"ZIP: Total number of files in {os.path.basename(dump_path)}: {len(archive.namelist())}") - for csv_name in archive.namelist(): - if csv_name.endswith('.csv'): - with archive.open(csv_name) as csv_file: - _p_csvfile(csv_file, csv_name, rconn_db_br, rconn_db_ra, rconn_db_metadata) - - elif dump_path.endswith(".tar.gz") or dump_path.endswith(".tgz"): - # Handle single TAR.GZ file - with tarfile.open(dump_path, 'r:gz') as archive: - logger.info(f"TAR.GZ: Total number of files in {os.path.basename(dump_path)}: {len(archive.getnames())}") - for csv_name in archive.getnames(): - if csv_name.endswith('.csv'): - csv_file = archive.extractfile(csv_name) - if csv_file: - _p_csvfile(csv_file, csv_name, rconn_db_br, rconn_db_ra, rconn_db_metadata) - - elif dump_path.endswith(".csv"): - # Handle single CSV file - csv_name = os.path.basename(dump_path) - logger.info(f"CSV: Processing direct CSV file: {csv_name}") - with open(dump_path, 'rb') as csv_file: - _p_csvfile(csv_file,csv_name, rconn_db_br, rconn_db_ra, rconn_db_metadata) - else: - logger.warning(f"Unsupported file type: {dump_path}") - - # Check if dump_path is a directory - elif os.path.isdir(dump_path): - # Directory contains only CSV files - for filename in os.listdir(dump_path): - filepath = os.path.join(dump_path, filename) - # Skip if it's not a file - if not os.path.isfile(filepath): - continue - - if filename.endswith(".csv"): - logger.info(f"CSV: Processing direct CSV file: {filename}") - with open(filepath, 'rb') as csv_file: - _p_csvfile(csv_file, filename, rconn_db_br, rconn_db_ra, rconn_db_metadata) - else: - logger.error(f"Path does not exist or is neither a file nor directory: {dump_path}") - - #print glob indexes to file - logger.info("Saving (in CSV) global indexes...") - with open('meta_br.csv', 'a+') as f: - write = csv.writer(f) - for any_id in br_index: - write.writerow([any_id,"; ".join(list(br_index[any_id]))]) - - with open('meta_ra.csv', 'a+') as f: - write = csv.writer(f) - for any_id in ra_index: - write.writerow([any_id,"; ".join(list(ra_index[any_id]))]) - - return (str(len(br_index)), str(len(ra_index))) - - -def main(): - global _config - - parser = argparse.ArgumentParser(description='Store the metadata of OpenCitations Meta in Redis') - parser.add_argument('--dump', type=str, required=True,help='The directory of CSVs or file (in ZIP or TAR.GZ) representing OpenCitations Meta dump') - #parser.add_argument('--db', type=str, required=True,help='The destination DB in redis. The specified DB is used to store : data, while DB+1 is used to store :{METADATA}') - #parser.add_argument('--port', type=str, required=False,help='The port of redis', default="6379") - - args = parser.parse_args() - logger = get_logger() - - logger.info("Start uploading data to Redis.") - - res = upload2redis( - dump_path = args.dump, - # Redis main conf - redishost = _config.get("redis", "host"), - redisport = _config.get("redis", "port"), - redisbatchsize = _config.get("redis", "batch_size"), - db_omid = _config.get("cnc", "db_omid"), - db_br = _config.get("cnc", "db_br"), - db_ra = _config.get("cnc", "db_ra"), - db_metadata = _config.get("INDEX", "db") - ) - - logger.info("A total of unique "+str(res[0])+" BR OMIDs and "+str(res[1])+" RA OMIDs have been found and added to Redis.") - -if __name__ == "__main__": - main() From a328ed4f88836fad023fabf56dc42982edd505c6 Mon Sep 17 00:00:00 2001 From: Hubert Krzywonos Date: Mon, 14 Sep 2026 12:36:54 +0200 Subject: [PATCH 4/4] configurable dump location, missing main function calls, specified encoding for loading CSV files, dump_index.process_pair() fix for non-fork-spawned global variables --- config.ini | 3 +++ oc_index/scripts/cnc.py | 9 ++++++--- oc_index/scripts/dump_index.py | 21 +++++++++++++-------- oc_index/scripts/meta2redis.py | 2 +- 4 files changed, 23 insertions(+), 12 deletions(-) diff --git a/config.ini b/config.ini index 8a04514..228d36e 100755 --- a/config.ini +++ b/config.ini @@ -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 diff --git a/oc_index/scripts/cnc.py b/oc_index/scripts/cnc.py index aeb3880..2ea29a0 100755 --- a/oc_index/scripts/cnc.py +++ b/oc_index/scripts/cnc.py @@ -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, ) @@ -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() \ No newline at end of file diff --git a/oc_index/scripts/dump_index.py b/oc_index/scripts/dump_index.py index cf6de90..f003a26 100644 --- a/oc_index/scripts/dump_index.py +++ b/oc_index/scripts/dump_index.py @@ -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): @@ -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**") @@ -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: @@ -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\"]}" @@ -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) @@ -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: @@ -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) @@ -343,3 +345,6 @@ def process_pair(pairs, pnum, br_meta, end_cursor = False): force = end_cursor, pnum = pnum ) + +if __name__ == "__main__": + main() \ No newline at end of file diff --git a/oc_index/scripts/meta2redis.py b/oc_index/scripts/meta2redis.py index 60ed4c8..24edbc5 100755 --- a/oc_index/scripts/meta2redis.py +++ b/oc_index/scripts/meta2redis.py @@ -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"])