From e5ae903aa2ccab1458808e8659cb05296c24dcef Mon Sep 17 00:00:00 2001 From: Romain Beaumont Date: Fri, 6 Jan 2023 03:09:05 +0100 Subject: [PATCH 1/7] add video platform does not actually work well --- cc2dataset/main.py | 18 ++++++++++++++++++ examples/single_warc_example.py | 2 +- 2 files changed, 19 insertions(+), 1 deletion(-) diff --git a/cc2dataset/main.py b/cc2dataset/main.py index 0d16574..59eaa4f 100644 --- a/cc2dataset/main.py +++ b/cc2dataset/main.py @@ -18,7 +18,21 @@ from .spark_session_builder import build_spark_session from io import BytesIO from urllib.parse import urljoin +from yt_dlp.extractor import gen_extractor_classes, GenericIE + +def valid_video_platform_link(link): + if "amazon" in link.get("url", "") or "drive" in link.get("url", "") or "twitter" in link.get("url", "") or "pinterest" in link.get("url", "") or "youtube" in link.get("url", "") or "instagram" in link.get("url", ""): + return False + for ie in gen_extractor_classes(): + if ie != GenericIE and ie.suitable(link.get("url", "")): + return True + return False + +def extract_video_platform_from_links(links): + #links = links[:100] + filtered_links = [{"url": link["url"], "alt": link.get("text", "")} for link in links if valid_video_platform_link(link)] + return filtered_links def valid_video_link(link): valid_video = any( @@ -127,6 +141,8 @@ def extract_documents_from_links(links, document_type): return extract_text_from_links(links) elif document_type == "video": return extract_video_from_links(links) + elif document_type == "video_platform": + return extract_video_platform_from_links(links) else: raise ValueError(f"Unknown document type {document_type}") @@ -175,6 +191,8 @@ def extract_documents_from_wat(stream, document_type): link["cc_filename"] = cc_filename link["page_url"] = page_url all_links.extend(filtered_links) + if len(all_links) > 100: + return all_links except Exception as e: # pylint: disable=broad-except logger.info(e) logger.info("A shard failed to parse") diff --git a/examples/single_warc_example.py b/examples/single_warc_example.py index 121ed3c..a247810 100644 --- a/examples/single_warc_example.py +++ b/examples/single_warc_example.py @@ -10,7 +10,7 @@ else: url = "https://data.commoncrawl.org/" + wat - results = process_wat(url, "image") + results = process_wat(url, "video_platform") df = pd.DataFrame(results, columns=["uid", "url", "alt", "cc_filename", "page_url"]) df.to_parquet(os.getcwd() + "/output.parquet") print(df) From 894a261ea149876d07b9be1c99e4e242813ff034 Mon Sep 17 00:00:00 2001 From: Romain Beaumont Date: Wed, 15 Nov 2023 22:56:06 +0100 Subject: [PATCH 2/7] wip --- cc2dataset/main.py | 69 ++++++++++++++++++++++++++-------------- examples/run_on_spark.py | 12 +++---- 2 files changed, 52 insertions(+), 29 deletions(-) diff --git a/cc2dataset/main.py b/cc2dataset/main.py index 59eaa4f..cb6bb4a 100644 --- a/cc2dataset/main.py +++ b/cc2dataset/main.py @@ -12,6 +12,7 @@ from pyspark import SparkContext from pyspark.sql.functions import rand from pyspark.sql import SparkSession +from bs4 import BeautifulSoup import random import math import time @@ -19,18 +20,32 @@ from io import BytesIO from urllib.parse import urljoin from yt_dlp.extractor import gen_extractor_classes, GenericIE - -def valid_video_platform_link(link): - if "amazon" in link.get("url", "") or "drive" in link.get("url", "") or "twitter" in link.get("url", "") or "pinterest" in link.get("url", "") or "youtube" in link.get("url", "") or "instagram" in link.get("url", ""): - return False - for ie in gen_extractor_classes(): - if ie != GenericIE and ie.suitable(link.get("url", "")): - return True +import re +def is_youtube_video(url): + if re.match('^https?://(www.)?youtube.com/watch\?v=.+$', url): + return True + if re.match('^https?://(www.)?youtube.com/v/.+$', url): + return True + if re.match('^https?://(www.)?youtube.com/embed/.+$', url): + return True + if re.match('^https?://(www.)?youtu.be/.+$', url): + return True + + return False + + +def is_bilibili_video(url): + if re.match("https?://(?:www\.)?bilibili\.com/(?:video/|festival/\w+\?(?:[^#]*&)?bvid=)[aAbB][vV](?P[^/?#&]+)",url): + return True + return False +def valid_video_platform_link(link): + return is_bilibili_video(link.get("url", "")) + + def extract_video_platform_from_links(links): - #links = links[:100] filtered_links = [{"url": link["url"], "alt": link.get("text", "")} for link in links if valid_video_platform_link(link)] return filtered_links @@ -205,16 +220,17 @@ def process_wat(path, document_type): """Process a single wat file""" begin_read = timer() with fsspec.open(path, "rb") as f: - for i in range(10): + retries = 1000 + for i in range(retries): try: tf = BytesIO(f.read()) break except Exception as ex: # pylint: disable=broad-except - if i == 9: - logger.info("failed 10 times, skipping ", path) + if i == retries-1: + logger.info(f"failed {retries} times, skipping ", path) return logger.info(ex) - logger.info(f"retrying reading {i}/10") + logger.info(f"retrying reading {i}/{retries}") time.sleep(1) for e in extract_documents_from_wat(tf, document_type): @@ -233,22 +249,29 @@ def get_cc_wat_links(source_cc_protocol): elif source_cc_protocol == "http": fs, p = fsspec.core.url_to_fs("https://commoncrawl.org/the-data/get-started/") a = fs.open(p).read() - l = a.splitlines() - l = [e.decode("utf8").replace("[WARC] ", "") for e in l] - l = [e for e in l if "
  • s3://commoncrawl/crawl-data/" in e] - l = [ - e.split(" ")[0].replace("
  • s3://commoncrawl/", "https://data.commoncrawl.org/").replace("", "") - for e in l - ] - l = [(e + "/wat.paths.gz").replace("//wat", "/wat") for e in l] - return l + soup = BeautifulSoup(a, 'html.parser') + h6_content = [e.text for e in soup.find_all('h6')][:-3] + h6_content= [e for e in h6_content if e != "CC-MAIN-2013-20" ] + results = [f"https://data.commoncrawl.org/crawl-data/{e}/wat.paths.gz" for e in h6_content] + return results else: raise ValueError(f"Unknown protocol {source_cc_protocol}") def read_wat_index_file(wat_index): - with fsspec.open(wat_index, "rb", compression="gzip") as f: - wats = [a.decode("utf8").strip() for a in f.readlines()] + retries = 1000 + for i in range(retries): + try: + with fsspec.open(wat_index, "rb", compression="gzip") as f: + wats = [a.decode("utf8").strip() for a in f.readlines()] + break + except Exception as ex: # pylint: disable=broad-except + if i == retries-1: + logger.info(f"failed {retries} times, skipping ", wat_index) + return + logger.info(ex) + logger.info(f"retrying reading {i}/{retries}") + time.sleep(1) return wats diff --git a/examples/run_on_spark.py b/examples/run_on_spark.py index f6720ad..d3e7a13 100644 --- a/examples/run_on_spark.py +++ b/examples/run_on_spark.py @@ -4,11 +4,11 @@ if __name__ == "__main__": # if you have a slurm cluster, refer to https://gist.github.com/rom1504/67ada3dedbecc113ae2dbdfd9c642d83 to start a spark cluster there cc2dataset( - "s3a://s-laion/cc-proc-test/tmpp", - wat_index_count=None, + "/tmp/tmp_output", wat_count=1000, - master="spark://cpu128-dy-r6i-32xlarge-27:7077", - num_cores=128, - mem_gb=256, - multipart=2, + master="local", + num_cores=32, + mem_gb=8, + document_type="video_platform", + source_cc_protocol="http" ) From 7bde88c76b42f0ed1746cb4b9e378e256fe085e7 Mon Sep 17 00:00:00 2001 From: Romain Beaumont Date: Thu, 16 Nov 2023 00:28:54 +0100 Subject: [PATCH 3/7] tmp2 --- cc2dataset/main.py | 35 +++++++++++++++++++++++++++++++---- examples/run_on_spark.py | 3 ++- 2 files changed, 33 insertions(+), 5 deletions(-) diff --git a/cc2dataset/main.py b/cc2dataset/main.py index cb6bb4a..7c3ad51 100644 --- a/cc2dataset/main.py +++ b/cc2dataset/main.py @@ -41,9 +41,38 @@ def is_bilibili_video(url): return False -def valid_video_platform_link(link): +def valid_video_platform_link_(link): return is_bilibili_video(link.get("url", "")) +import yt_dlp + +generic_extractors = [yt_dlp.extractor.generic.GenericIE, yt_dlp.extractor.lazy_extractors.GenericIE] + +FILTERED_EXTRACTORS = {ie.IE_NAME:ie for ie in yt_dlp.list_extractor_classes() + if ie not in generic_extractors + and "porn" not in ie.IE_NAME.lower() + and "adult" not in ie.IE_NAME.lower() + and "xxx" not in ie.IE_NAME.lower() + and "xvideos" not in ie.IE_NAME.lower() + and "xhamster" not in ie.IE_NAME.lower() + and "redtube" not in ie.IE_NAME.lower() + and "xtube" not in ie.IE_NAME.lower() + and "xstream" not in ie.IE_NAME.lower() + and "xfileshare" not in ie.IE_NAME.lower() + and "sex" not in ie.IE_NAME.lower() + } + + +# print(FILTERED_EXTRACTORS.keys()) +# print(len(FILTERED_EXTRACTORS.keys())) + +def is_link_valid(link, extractors): + """Check if link is valid given a list of extractors.""" + return any([ie.suitable(link) for ie in extractors]) + +def valid_video_platform_link(link): + """Check if link is a valid video platform link.""" + return is_link_valid(link.get("url", ""), FILTERED_EXTRACTORS.values()) def extract_video_platform_from_links(links): filtered_links = [{"url": link["url"], "alt": link.get("text", "")} for link in links if valid_video_platform_link(link)] @@ -206,8 +235,6 @@ def extract_documents_from_wat(stream, document_type): link["cc_filename"] = cc_filename link["page_url"] = page_url all_links.extend(filtered_links) - if len(all_links) > 100: - return all_links except Exception as e: # pylint: disable=broad-except logger.info(e) logger.info("A shard failed to parse") @@ -251,7 +278,7 @@ def get_cc_wat_links(source_cc_protocol): a = fs.open(p).read() soup = BeautifulSoup(a, 'html.parser') h6_content = [e.text for e in soup.find_all('h6')][:-3] - h6_content= [e for e in h6_content if e != "CC-MAIN-2013-20" ] + h6_content= [e for e in h6_content ] results = [f"https://data.commoncrawl.org/crawl-data/{e}/wat.paths.gz" for e in h6_content] return results else: diff --git a/examples/run_on_spark.py b/examples/run_on_spark.py index d3e7a13..c9f329b 100644 --- a/examples/run_on_spark.py +++ b/examples/run_on_spark.py @@ -5,7 +5,8 @@ # if you have a slurm cluster, refer to https://gist.github.com/rom1504/67ada3dedbecc113ae2dbdfd9c642d83 to start a spark cluster there cc2dataset( "/tmp/tmp_output", - wat_count=1000, + wat_index_count=None, + wat_count=10, master="local", num_cores=32, mem_gb=8, From d0d8037d0a4275802a6bed0cd96c4dc5d0e378a0 Mon Sep 17 00:00:00 2001 From: Romain Beaumont Date: Fri, 17 Nov 2023 00:17:35 +0100 Subject: [PATCH 4/7] faster --- cc2dataset/main.py | 54 +++++++++++++++++++++++++++++++++++++++++----- 1 file changed, 49 insertions(+), 5 deletions(-) diff --git a/cc2dataset/main.py b/cc2dataset/main.py index 7c3ad51..cb464f3 100644 --- a/cc2dataset/main.py +++ b/cc2dataset/main.py @@ -20,6 +20,8 @@ from io import BytesIO from urllib.parse import urljoin from yt_dlp.extractor import gen_extractor_classes, GenericIE +from urllib.parse import urlparse +import traceback import re def is_youtube_video(url): @@ -62,17 +64,56 @@ def valid_video_platform_link_(link): and "sex" not in ie.IE_NAME.lower() } +def extract_test(extractor): + tests = [] + if hasattr(extractor, "_TEST") and extractor._TEST is not None: + tests = [extractor._TEST["url"]] + elif hasattr(extractor, "_TESTS") and extractor._TESTS is not None: + tests = [x["url"] for x in extractor._TESTS] + return tests + +def normalize_domain(domain): + domain = domain.lower() + if domain.startswith("www."): + domain = domain[4:] + return domain + +def extract_domain(url): + try: + parsed_url = urlparse(url) + domain = parsed_url.netloc + return normalize_domain(domain) + except: + return "" + +DOMAIN_DICT = {} + +for extractor in FILTERED_EXTRACTORS.values(): + for url in extract_test(extractor): + domain = extract_domain(url) + if domain in DOMAIN_DICT: + DOMAIN_DICT[domain] = DOMAIN_DICT[domain] + [extractor] + else: + DOMAIN_DICT[domain] = [extractor] -# print(FILTERED_EXTRACTORS.keys()) -# print(len(FILTERED_EXTRACTORS.keys())) +def is_link_suitable(link, extractors): + """Check if link is valid given an extractor.""" + try: + return any([ie.suitable(link) for ie in extractors]) + except: + return False -def is_link_valid(link, extractors): +def is_link_valid(link, domain_dict): """Check if link is valid given a list of extractors.""" - return any([ie.suitable(link) for ie in extractors]) + is_valid = False + domain = extract_domain(link) + if domain in domain_dict: + is_valid = is_link_suitable(link, domain_dict[domain]) + return is_valid def valid_video_platform_link(link): """Check if link is a valid video platform link.""" - return is_link_valid(link.get("url", ""), FILTERED_EXTRACTORS.values()) + return is_link_valid(link.get("url", ""), DOMAIN_DICT) def extract_video_platform_from_links(links): filtered_links = [{"url": link["url"], "alt": link.get("text", "")} for link in links if valid_video_platform_link(link)] @@ -235,7 +276,10 @@ def extract_documents_from_wat(stream, document_type): link["cc_filename"] = cc_filename link["page_url"] = page_url all_links.extend(filtered_links) + #if len(all_links) > 1000: + # return all_links except Exception as e: # pylint: disable=broad-except + traceback.print_exc() logger.info(e) logger.info("A shard failed to parse") return [] From 9af33835ad43efd1e83ef859ca44c76adb40019e Mon Sep 17 00:00:00 2001 From: Romain Beaumont Date: Thu, 7 Dec 2023 20:58:30 +0100 Subject: [PATCH 5/7] improve filters --- cc2dataset/main.py | 68 +++++++++++++++++++++++++--------------- examples/run_on_spark.py | 4 +-- 2 files changed, 45 insertions(+), 27 deletions(-) diff --git a/cc2dataset/main.py b/cc2dataset/main.py index cb464f3..8e4283c 100644 --- a/cc2dataset/main.py +++ b/cc2dataset/main.py @@ -47,21 +47,31 @@ def valid_video_platform_link_(link): return is_bilibili_video(link.get("url", "")) import yt_dlp +import unicodedata generic_extractors = [yt_dlp.extractor.generic.GenericIE, yt_dlp.extractor.lazy_extractors.GenericIE] - -FILTERED_EXTRACTORS = {ie.IE_NAME:ie for ie in yt_dlp.list_extractor_classes() - if ie not in generic_extractors - and "porn" not in ie.IE_NAME.lower() - and "adult" not in ie.IE_NAME.lower() - and "xxx" not in ie.IE_NAME.lower() - and "xvideos" not in ie.IE_NAME.lower() - and "xhamster" not in ie.IE_NAME.lower() - and "redtube" not in ie.IE_NAME.lower() - and "xtube" not in ie.IE_NAME.lower() - and "xstream" not in ie.IE_NAME.lower() - and "xfileshare" not in ie.IE_NAME.lower() - and "sex" not in ie.IE_NAME.lower() +porn_patterns = ["porn", "adult", "xxx", "xvideos", "xhamster", "redtube", "xtube", "xstream", "xfileshare", "sex"] +playlist_patterns = ["Playlist", "Category", "User"] +domain_patterns = ["twitter", "instagram", "facebook", "player.zype", "imgur", "flickr"] +youtube_whitelist = ["YoutubeIE", "YoutubeYtBeIE", "YoutubeClipIE"] +dailymotion_whitelist = ["DailymotionIE"] + +def get_class_name(ie): + return str(ie).split('.')[-1].split("'")[0] + +def substrings_not_in_string(s, subs): + not_in_string = [ss for ss in subs if ss in s] + return not not_in_string + +def whitlist_extractors(ie, main_name, extractor_whitelist): + return not main_name in ie.IE_NAME.lower() or get_class_name(ie) in youtube_whitelist + +FILTERED_EXTRACTORS = {ie.IE_NAME:ie for ie in yt_dlp.list_extractor_classes() + if ie not in generic_extractors + and substrings_not_in_string(ie.IE_NAME.lower(), porn_patterns) + and whitlist_extractors(ie, "youtube", youtube_whitelist) + and substrings_not_in_string(get_class_name(ie), playlist_patterns) + and substrings_not_in_string(get_class_name(ie).lower(), domain_patterns) } def extract_test(extractor): @@ -73,28 +83,36 @@ def extract_test(extractor): return tests def normalize_domain(domain): - domain = domain.lower() - if domain.startswith("www."): + domain = domain.lower() + if domain.startswith("www."): domain = domain[4:] - return domain + return domain + +def normalize_url(url): + normalized_url = unicodedata.normalize('NFKC', url) + return normalized_url def extract_domain(url): try: - parsed_url = urlparse(url) + parsed_url = urlparse(normalize_url(url)) domain = parsed_url.netloc return normalize_domain(domain) - except: + except Exception as e: return "" -DOMAIN_DICT = {} + +DOMAIN_IES_DICT = {} for extractor in FILTERED_EXTRACTORS.values(): for url in extract_test(extractor): domain = extract_domain(url) - if domain in DOMAIN_DICT: - DOMAIN_DICT[domain] = DOMAIN_DICT[domain] + [extractor] - else: - DOMAIN_DICT[domain] = [extractor] + if domain in DOMAIN_IES_DICT: + if extractor not in DOMAIN_IES_DICT[domain]: + DOMAIN_IES_DICT[domain] = DOMAIN_IES_DICT[domain] + [extractor] + else: + DOMAIN_IES_DICT[domain] = [extractor] + + def is_link_suitable(link, extractors): """Check if link is valid given an extractor.""" @@ -108,12 +126,12 @@ def is_link_valid(link, domain_dict): is_valid = False domain = extract_domain(link) if domain in domain_dict: - is_valid = is_link_suitable(link, domain_dict[domain]) + is_valid = is_link_suitable(link, domain_dict[domain]) return is_valid def valid_video_platform_link(link): """Check if link is a valid video platform link.""" - return is_link_valid(link.get("url", ""), DOMAIN_DICT) + return is_link_valid(link.get("url", ""), DOMAIN_IES_DICT) def extract_video_platform_from_links(links): filtered_links = [{"url": link["url"], "alt": link.get("text", "")} for link in links if valid_video_platform_link(link)] diff --git a/examples/run_on_spark.py b/examples/run_on_spark.py index c9f329b..5b25056 100644 --- a/examples/run_on_spark.py +++ b/examples/run_on_spark.py @@ -6,9 +6,9 @@ cc2dataset( "/tmp/tmp_output", wat_index_count=None, - wat_count=10, + wat_count=1000, master="local", - num_cores=32, + num_cores=16, mem_gb=8, document_type="video_platform", source_cc_protocol="http" From 076358e836ed86ab780249e2ae22c53ebf05a454 Mon Sep 17 00:00:00 2001 From: Romain Beaumont Date: Fri, 8 Dec 2023 23:00:05 +0100 Subject: [PATCH 6/7] remove empty domain --- cc2dataset/main.py | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/cc2dataset/main.py b/cc2dataset/main.py index 8e4283c..81656d2 100644 --- a/cc2dataset/main.py +++ b/cc2dataset/main.py @@ -106,6 +106,8 @@ def extract_domain(url): for extractor in FILTERED_EXTRACTORS.values(): for url in extract_test(extractor): domain = extract_domain(url) + if domain == "": + continue if domain in DOMAIN_IES_DICT: if extractor not in DOMAIN_IES_DICT[domain]: DOMAIN_IES_DICT[domain] = DOMAIN_IES_DICT[domain] + [extractor] @@ -113,7 +115,6 @@ def extract_domain(url): DOMAIN_IES_DICT[domain] = [extractor] - def is_link_suitable(link, extractors): """Check if link is valid given an extractor.""" try: From a30d37970c675159797e09791c017a836571f4e0 Mon Sep 17 00:00:00 2001 From: Romain Beaumont Date: Sat, 9 Dec 2023 15:43:52 +0100 Subject: [PATCH 7/7] Make wat retries count configurable. --- cc2dataset/main.py | 17 +++++++++-------- 1 file changed, 9 insertions(+), 8 deletions(-) diff --git a/cc2dataset/main.py b/cc2dataset/main.py index 81656d2..ed5cb17 100644 --- a/cc2dataset/main.py +++ b/cc2dataset/main.py @@ -348,8 +348,8 @@ def get_cc_wat_links(source_cc_protocol): raise ValueError(f"Unknown protocol {source_cc_protocol}") -def read_wat_index_file(wat_index): - retries = 1000 +def read_wat_index_file(wat_index, wat_read_retries): + retries = wat_read_retries for i in range(retries): try: with fsspec.open(wat_index, "rb", compression="gzip") as f: @@ -397,7 +397,7 @@ def deduplicate_repartition_count(df, output_path, wat_count, spark, shuffle=Fal logger.info(f"Size: {df.count()}") -def process_one_part(output_path, wat_index_files, build_spark, shuffle, document_type, source_cc_protocol): +def process_one_part(output_path, wat_index_files, build_spark, shuffle, document_type, source_cc_protocol, wat_read_retries): """Process one part""" spark = build_spark() sc = SparkContext.getOrCreate() @@ -410,7 +410,7 @@ def process_one_part(output_path, wat_index_files, build_spark, shuffle, documen def extract(x): x = list(x) - yield from process_wat(prefix + x[0], document_type) + yield from process_wat(prefix + x[0], document_type, wat_read_retries) output = wat_rdd.mapPartitions(extract) df = output.toDF(["uid", "url", "alt", "cc_filename", "page_url"]) @@ -428,7 +428,7 @@ def get_last_successful_part(output_path): def process_multi_part( - output_path, wat_index_files, build_spark, multipart, shuffle, resume, document_type, source_cc_protocol + output_path, wat_index_files, build_spark, multipart, shuffle, resume, document_type, source_cc_protocol, wat_read_retries ): """Process multi part""" if resume: @@ -445,7 +445,7 @@ def process_multi_part( part_path = f"{output_path}/part_{i}" part_paths.append(part_path) logger.info(f"Processing part {i} from {start} to {end} into {part_path}") - process_one_part(part_path, wat_index_files[start:end], build_spark, False, document_type, source_cc_protocol) + process_one_part(part_path, wat_index_files[start:end], build_spark, False, document_type, source_cc_protocol, wat_read_retries) spark = build_spark() logger.info("Merging parts") @@ -477,6 +477,7 @@ def cc2dataset( spark_builder=None, document_type="image", source_cc_protocol="s3", + wat_read_retries=1000, ): """Convert common crawl to image caption set""" @@ -511,10 +512,10 @@ def build_spark(): wat_index_files = f.read().splitlines() if multipart is None: - process_one_part(output_path, wat_index_files, build_spark, shuffle, document_type, source_cc_protocol) + process_one_part(output_path, wat_index_files, build_spark, shuffle, document_type, source_cc_protocol, wat_read_retries) else: process_multi_part( - output_path, wat_index_files, build_spark, multipart, shuffle, resume, document_type, source_cc_protocol + output_path, wat_index_files, build_spark, multipart, shuffle, resume, document_type, source_cc_protocol, wat_read_retries )