Palantir Foundry Sidecar Transform超时问题排查求助
问题背景
我正在Palantir Foundry中搭建一个基于Selenium的网页爬取POC项目,目标是爬取数十万量级的网站。我清楚这个方案可能不是最优解,但希望找到可行路径,也欢迎Palantir相关人员指出不可行性。静态内容用Beautiful Soup在常规Transform中运行正常,但因为需要爬取动态生成的内容,所以考虑用容器+Sidecar Transform运行Selenium,如有更优方案也请建议。
当前我的Sidecar Transform运行20-30分钟后超时失败,错误信息如下:
[module version: 1.1132.0]
Spark module 'ri.spark-module-manager.main.spark-module.e3afe96c-4d51-44a4-a687-1174dfba2fb4' died while job 'ri.foundry.main.job.f00598ef-e11e-43c4-82bf-ac63d057294a' was using it. (ExitReason: MODULE_UNREACHABLE)
Module exit details: Module became unreachable after registration. This likely indicates the module has died. Module became unreachable for an unknown reason.
代码实现
Sidecar Transform代码
from transforms.api import transform, Input, Output, configure from transforms.sidecar import sidecar, Volume from myproject.datasets.utils import copy_files_to_shared_directory, copy_output_files from myproject.datasets.utils import copy_start_flag, wait_for_done_flag, copy_close_flag, launch_udf_once from transforms.external.systems import use_external_systems, EgressPolicy, Credential @use_external_systems( egress=EgressPolicy('{POLICY RID}') ) @configure(["NUM_EXECUTORS_64", 'EXECUTOR_MEMORY_LARGE', 'EXECUTOR_MEMORY_OVERHEAD_LARGE', 'DRIVER_MEMORY_EXTRA_EXTRA_LARGE', 'DRIVER_MEMORY_OVERHEAD_EXTRA_LARGE' ]) @sidecar(image='{PACKAGE_NAME}', tag='0.3', volumes=[Volume("shared")]) @transform( output=Output("OUTPUT"), ) def compute(ctx, output, egress): def user_defined_function(row): # Copy files from source to shared directory. # copy_files_to_shared_directory(source) # Send the start flag so the container knows it has all the input files copy_start_flag() # Iterate till the stop flag is written or we hit the max time limit wait_for_done_flag() # Copy out output files from the container to an output dataset output_fnames = [ "start_flag", # "outfile.csv", "logfile", "done_flag", ] copy_output_files(output, output_fnames) # Write the close flag so the container knows you have extracted the data copy_close_flag() # The user defined function must return something return (row.ExecutionID, "success") # This spawns one task, which maps to one executor, and launches one "sidecar container" launch_udf_once(ctx, user_defined_function)
Dockerfile
FROM --platform=linux/amd64 python:3.9-buster RUN mkdir /code # Keeps Python from generating .pyc files in the container ENV PYTHONDONTWRITEBYTECODE=1 # Turns off buffering for easier container logging ENV PYTHONUNBUFFERED=1 # please review all the latest versions here: # https://googlechromelabs.github.io/chrome-for-testing/ ENV CHROMEDRIVER_VERSION=123.0.6312.122 ### install chrome # https://storage.googleapis.com/chrome-for-testing-public/123.0.6312.122/linux64/chrome-linux64.zip RUN apt-get update && apt-get install -y wget && apt-get install -y zip # RUN wget -q https://dl.google.com/linux/direct/google-chrome-stable_current_amd64.deb # RUN apt-get install -y ./google-chrome-stable_current_amd64.deb COPY google-chrome-stable_current_amd64.deb . RUN apt-get install -y ./google-chrome-stable_current_amd64.deb ### install chromedriver # RUN wget https://storage.googleapis.com/chrome-for-testing-public/123.0.6312.122/linux64/chromedriver-linux64.zip \ # && unzip chromedriver-linux64.zip && rm -dfr chromedriver_linux64.zip \ # && mv /chromedriver-linux64/chromedriver /usr/bin/chromedriver \ # && chmod +x /usr/bin/chromedriver COPY chromedriver-linux64.zip . RUN unzip chromedriver-linux64.zip && rm -dfr chromedriver_linux64.zip \ && mv /chromedriver-linux64/chromedriver /usr/bin/chromedriver \ && chmod +x /usr/bin/chromedriver # set display port to avoid crash ENV DISPLAY=:99 # install selenium RUN pip install selenium==4.3.0 ADD entrypoint.py /usr/bin/ ADD scraper.py /usr/bin/ RUN chmod +x /usr/bin/ RUN mkdir -p /opt/palantir/sidecars/shared-volumes/shared/ RUN chown 5001 /opt/palantir/sidecars/shared-volumes/shared/ ENV SHARED_DIR=/opt/palantir/sidecars/shared-volumes/shared USER 5001 CMD ["/usr/bin/entrypoint.py"] ENTRYPOINT ["python"]
entrypoint.py
import os import time import subprocess from datetime import datetime import argparse def run_process(): "Define a function for running commands and capturing stdout line by line" p = subprocess.Popen(["python", "/usr/bin/scraper.py"], stdout=subprocess.PIPE, stderr=subprocess.STDOUT) out, err = p.communicate() return (p.returncode, out, err) # debug ''' item = run_process() my_string = f"{datetime.utcnow().isoformat()}: {item}" print(my_string) ''' start_flag_fname = "/opt/palantir/sidecars/shared-volumes/shared/start_flag" done_flag_fname = "/opt/palantir/sidecars/shared-volumes/shared/done_flag" close_flag_fname = "/opt/palantir/sidecars/shared-volumes/shared/close_flag" # Wait for start flag print(f"{datetime.utcnow().isoformat()}: waiting for start flag") while not os.path.exists(start_flag_fname): time.sleep(1) print(f"{datetime.utcnow().isoformat()}: start flag detected") # Execute model, logging output to file with open("/opt/palantir/sidecars/shared-volumes/shared/logfile", "w") as logfile: item = run_process() my_string = f"{datetime.utcnow().isoformat()}: {item}" print(my_string) logfile.write(my_string) logfile.flush() print(f"{datetime.utcnow().isoformat()}: execution finished writing output file") # Write out the done flag open(done_flag_fname, "w") print(f"{datetime.utcnow().isoformat()}: done flag file written") # Wait for close flag before allowing the script to finish while not os.path.exists(close_flag_fname): time.sleep(1) print(f"{datetime.utcnow().isoformat()}: close flag detected. shutting down")
scraper.py
from selenium import webdriver from selenium.webdriver.chrome.options import Options # Define options for running the chromedriver chrome_options = Options() chrome_options.add_argument("--no-sandbox") chrome_options.add_argument("--headless") chrome_options.add_argument("--disable-dev-shm-usage") # Initialize a new chrome driver instance driver = webdriver.Chrome(options=chrome_options) driver.get('{WEBSITE}') print(driver.page_source) driver.quit()
请求协助
恳请协助排查超时原因及给出调试方向。
内容的提问来源于stack exchange,提问作者tessa

