Spark DataFrame提取特定字符串并按规则拆分列的需求
Spark DataFrame id1列拆分实现
需求说明
给定如下Spark DataFrame:
df = [{'id': 1, 'id1': '859A;'}, {'id': 2, 'id1': '209A/229A/509A;'}, {'id': 3, 'id1': '(105A/111A/121A/131A/201A/205A/211A/221A/231A/509A/801A/805A/811A/821A)+TZ+-494;'}, {'id': 4, 'id1': '111A/114A/121A/131A/201A/211A/221A/231A/651A+-(Y05/U17)/801A/804A/821A;'}, {'id': 5, 'id1': '(651A/851A)+U17/861A;'}, ] df = spark.createDataFrame(df)
需要将id1列拆分为两列:
- newcolumn1:提取所有以"A"结尾的字符串,按规则整理:同一数字开头且以"A"结尾的字符串用逗号分隔为一组,不同组之间用"/"分隔,保留原有的括号结构(无括号则添加括号包裹);
- newcolumn2:提取剩余的非"A"结尾的字符串内容,无剩余则为空字符串。
解决方案
通过自定义UDF实现复杂的字符串解析和分组逻辑,具体代码如下:
import re from pyspark.sql.functions import udf from pyspark.sql.types import StructType, StructField, StringType # 定义UDF返回值的Schema result_schema = StructType([ StructField("newcolumn1", StringType(), nullable=True), StructField("newcolumn2", StringType(), nullable=True) ]) def group_a_strings(a_list): """将A结尾字符串按开头数字分组,同组用逗号分隔,不同组用/分隔""" groups = {} for s in a_list: # 提取字符串开头的数字作为分组键 key = s[0] if key not in groups: groups[key] = [] groups[key].append(s) # 按键排序后拼接分组内容 sorted_keys = sorted(groups.keys()) return '/'.join([','.join(groups[k]) for k in sorted_keys]) def process_id1(s): # 移除结尾的分号 s_clean = s.rstrip(';') newcolumn1 = "" newcolumn2 = "" # 分离括号内的内容和外部内容 bracket_match = re.match(r'^\((.*?)\)(.*)$', s_clean) has_bracket = False bracket_a_strings = [] non_bracket_content = s_clean if bracket_match: has_bracket = True bracket_inner = bracket_match.group(1) non_bracket_content = bracket_match.group(2) # 提取括号内所有以A结尾的字符串 bracket_a_strings = re.findall(r'\b\w+A\b', bracket_inner) # 提取非括号部分的A结尾字符串和非A内容 non_a_parts = [] non_bracket_a_strings = [] # 按特殊符号分割,区分A结尾和非A部分 segments = re.split(r'(\+|\-|\/)', non_bracket_content) for seg in segments: if re.fullmatch(r'\b\w+A\b', seg): non_bracket_a_strings.append(seg) else: if seg.strip(): non_a_parts.append(seg) # 处理newcolumn1的内容 if has_bracket: # 分别处理括号内和括号外的A字符串分组,合并到括号内 bracket_grouped = group_a_strings(bracket_a_strings) non_bracket_grouped = group_a_strings(non_bracket_a_strings) combined_grouped = f"{bracket_grouped},{non_bracket_grouped}" if non_bracket_grouped else bracket_grouped newcolumn1 = f"({combined_grouped})" else: # 无括号则直接分组后包裹括号 all_a_grouped = group_a_strings(bracket_a_strings + non_bracket_a_strings) newcolumn1 = f"({all_a_grouped})" if all_a_grouped else "" # 处理newcolumn2的内容,添加结尾分号 newcolumn2 = ''.join(non_a_parts) + ';' if non_a_parts else "" return (newcolumn1, newcolumn2) # 注册UDF process_id1_udf = udf(process_id1, result_schema) # 应用UDF生成新列 result_df = df.withColumn("split_result", process_id1_udf(df.id1)) \ .select( "id", "id1", "split_result.newcolumn1", "split_result.newcolumn2" ) # 查看结果 result_df.show(truncate=False)
验证结果
运行上述代码后,关键行的输出与需求示例一致:
- id=2:
newcolumn1:(209A,229A/509A)
newcolumn2: ``(空字符串) - id=3:
newcolumn1:(105A,111A,121A,131A/201A,205A,211A,221A,231A/509A/801A,805A,811A,821A)
newcolumn2:+TZ+-494; - id=5:
newcolumn1:(651A/851A,861A)
newcolumn2:+U17;
内容的提问来源于stack exchange,提问作者penchalaiah narakatla
相关产品推荐
相关产品推荐

