Pipeline nhanh
Trước khi bạn phân tích dữ liệu phức tạp hơn, quản lý muốn xem một ví dụ pipeline đơn giản gồm các bước cơ bản. Ở ví dụ này, bạn sẽ nạp một tệp dữ liệu, lọc một vài hàng, thêm một cột ID, rồi ghi ra dưới dạng dữ liệu JSON.
Ngữ cảnh spark đã được định nghĩa, và thư viện pyspark.sql.functions đã được đặt bí danh là F như thông lệ.
Bài tập này là một phần của khóa học
Làm sạch dữ liệu với PySpark
Hướng dẫn bài tập
- Nhập tệp
2015-departures.csv.gzvào một DataFrame. Lưu ý phần header đã được xác định sẵn. - Lọc DataFrame để chỉ giữ các chuyến bay có thời lượng lớn hơn 0 phút. Sử dụng chỉ số của cột, không dùng tên cột (hãy nhớ dùng
.printSchema()để xem tên cột/thứ tự cột). - Thêm một cột ID.
- Ghi tệp ra dưới dạng tài liệu JSON tên
output.json.
Bài tập tương tác thực hành trực tiếp
Hãy thử làm bài tập này bằng cách hoàn thành đoạn mã mẫu này.
# Import the data to a DataFrame
departures_df = spark.____(____, header=____)
# Remove any duration of 0
departures_df = departures_df.____(____)
# Add an ID column
departures_df = departures_df.____('id', ____)
# Write the file out to JSON format
____.write.____(____, mode='overwrite')