Showing posts with label pyspark. Show all posts
Showing posts with label pyspark. Show all posts

Wednesday, December 7, 2022

Conditional DataFrame column operations

Sử dụng when và otherwise

Chẳng hạn bạn thêm một cột 'tên_côt_gioi_tinh', các giá trị trong cột này sẽ là N nếu giá trị ở cột gioi_tinh tương ứng là 'Nam', còn lại thì sẽ là G.

df = df.withColumn('ten_cot_gioi_tinh', when(df.gioi_tinh == 'Nam', N)

                                                                .otherwise(G))

Thursday, November 24, 2022

Cleaning Data with PySpark

Pyspark dataframe column opertions: filter, select, withColumn.

Xử lý dữ liệu thô để sử dụng trong data processing pipeline, nếu dữ liệu không được làm sạch, thì sẽ gây ra những vấn đề sau đó liên quan tới hiệu năng và tổ chức luồng dữ liệu.

  • Định dạng kiểu dữ liệu, thay thế văn bản
  • Các chuyển đổi tính toán 
  • Loại bỏ dữ liệu "rác" và dữ liệu chưa hoàn thiện. 
Spark có những ưu điểm - đó là nâng cấp mở rộng được nguồn dữ liệu (scalable) và có bộ khung sườn  xử lý hiệu quả dữ liệu. 

Khi nhập file dữ liệu sử dụng pyspark thì có thêm một đối biến schema, spark schema là gì? 

spark schema là cấu trúc của DataFrame hoặc bộ dữ liệu, chúng ta định nghĩa cấu trúc này bằng cách sử dụng class StructTtype - là tổng hợp các StructField định nghĩa tên cột, kiểu dữ liệu của cột, cột nullable, và MetaData. 

Ví dụ dưới đây:

import pyspark.sql.types

yourSchema = StructType([

  StructField('ten', StringType(), True), 

  StructField('tuoi', IntegerType(), True),

  StructField('thanh pho', StringType(), True)

])

Giờ tiến nhành đọc file sử dụng schema như trên:

dan_cu = spark.read.format('csv').load(name='du_lieu_dan_cu.csv', schema=yourSchema)

Cách lấy - Get the last entry of the splits list

Sau khi tách string thành list thì bạn muốn lấy entry cuối, ở đây mình tạo một cột mới "Họ", còn cột_tách chứa list of strings họ và tên:

df = df.withColumn("Họ", df.splits.getItem(F.size('cột_tách') - 1))