在GCP Dataproc运行PySpark Job失败,请求排查代码问题
代码中的潜在问题及修复方案
你的PySpark任务失败大概率是由以下几个问题导致的:
1. HDFS路径错误
Dataproc集群内部访问HDFS时,无需指定Master节点实例名,直接使用集群默认的HDFS根路径即可。你当前的路径hdfs://{}/data/logs_example/*会导致namenode地址解析失败,正确路径应为hdfs:///data/logs_example/*(三个斜杠代表使用集群默认namenode)。
2. 日志分割逻辑脆弱,易触发索引越界
直接用split(" ")分割日志行非常不可靠:
- 日志日期字段通常带括号(如
[10/Oct/2000:13:55:36 -0700]),split后会被拆成多个元素 - URL或其他字段可能包含空格,导致分割后的数组长度与预期不符
- 一旦某行日志格式异常,访问
x[0]、x[3]等索引时会抛出IndexError,直接导致任务崩溃
3. 缺少异常处理逻辑
没有对格式错误的日志行做过滤,只要存在一行不符合预期的日志,整个任务就会失败。
修正后的代码
from pyspark.sql import SparkSession import re spark = SparkSession.builder \ .appName('spark_hdfs_to_hdfs') \ .getOrCreate() sc = spark.sparkContext sc.setLogLevel("WARN") # 使用集群默认HDFS路径 log_files_rdd = sc.textFile('hdfs:///data/logs_example/*') # 用正则表达式解析标准Nginx/Apache日志格式,避免空格分割的问题 LOG_PATTERN = r'^(\S+) \S+ \S+ \[([^\]]+)\] "(\S+) (\S+).*"' def parse_log_line(line): match = re.match(LOG_PATTERN, line) if match: return (match.group(1), match.group(2), match.group(3), match.group(4)) # 返回None过滤格式错误的行 return None # 过滤解析失败的日志行 parsed_rdd = log_files_rdd.map(parse_log_line).filter(lambda x: x is not None) columns = ["ip","date","method","url"] logs_df = parsed_rdd.toDF(columns) logs_df.createOrReplaceTempView('logs_df') sql = """ SELECT url, count(*) as count FROM logs_df WHERE url LIKE '%/article%' GROUP BY url """ article_count_df = spark.sql(sql) print(" ### Get only articles and blogs records ### ") article_count_df.show(5) # 显式停止SparkSession spark.stop()
内容的提问来源于stack exchange,提问作者olumide adeniyi
相关产品推荐
相关产品推荐

