Skip to main content

A package to load and preprocess JSON data using PySpark

Project description

PySpark JSON Loader

Overview

pyspark-json-loader is a Python package designed to facilitate loading and preprocessing JSON data using PySpark. It provides functions to start a Spark session, connect to a PostgreSQL database, preprocess data, and convert Spark DataFrames to Pandas DataFrames.

Installation

To install the package, run:

pip install pyspark-json-loader

Usage

Here is a detailed guide on how to use the functions provided by the pyspark-json-loader package.

Importing the Module

To use the functions in this package, you need to import them as follows:

from pyspark_json_loader import (
    load_json_file,
    start_spark_session,
    connect_to_db,
    preprocess_dataframe,
    convert_to_pandas,
    split_and_convert,
    add_lists,
    increment_index,
    index,
    multiply_lists,
    multiplyP_lists,
    sum_list,
    list_to_colon_separated_string
)

Functions

1. load_json_file(file_name)

Description: This function loads a JSON file and returns its contents as a Python dictionary.

Parameters:

  • file_name (str): The path to the JSON file.

Returns:

  • dict: The contents of the JSON file.

Example:

json_data = load_json_file('data.json')
print(json_data)

Output:

{
    "column1": "value1",
    "column2": "value2"
}

2. start_spark_session()

Description: This function starts a Spark session with the specified JAR file.

Returns:

  • SparkSession: The Spark session object.

Example:

spark = start_spark_session()
print(spark)

Output:

<pyspark.sql.session.SparkSession object at 0x...>

3. connect_to_db(spark, host_name, port_number, db_name, user, password, query, null_value=0)

Description: This function connects to a PostgreSQL database using the provided connection details and query.

Parameters:

  • spark (SparkSession): The Spark session object.
  • host_name (str): The hostname of the PostgreSQL server.
  • port_number (int): The port number of the PostgreSQL server.
  • db_name (str): The name of the database.
  • user (str): The username for the database.
  • password (str): The password for the database.
  • query (str): The SQL query to execute.
  • null_value (int, optional): The value to use for null values. Default is 0.

Returns:

  • DataFrame: The resulting DataFrame from the query.

Example:

df = connect_to_db(spark, 'localhost', 5432, 'mydatabase', 'user', 'password', 'SELECT * FROM mytable')
df.show()

Output:

+----+------+
| id | name |
+----+------+
| 1  | John |
| 2  | Jane |
+----+------+

4. preprocess_dataframe(df, json_data)

Description: This function preprocesses a DataFrame by filling null values based on the provided JSON data.

Parameters:

  • df (DataFrame): The Spark DataFrame to preprocess.
  • json_data (dict): The JSON data to use for preprocessing.

Returns:

  • DataFrame: The preprocessed DataFrame.

Example:

preprocessed_df = preprocess_dataframe(df, json_data)
preprocessed_df.show()

Output:

+---+------+
| id|  name|
+---+------+
| 1 |  John|
| 2 |  Jane|
+---+------+

5. convert_to_pandas(grouped_metrics_df)

Description: This function converts a Spark DataFrame to a Pandas DataFrame and adds a UUID column.

Parameters:

  • grouped_metrics_df (DataFrame): The Spark DataFrame to convert.

Returns:

  • DataFrame: The resulting Pandas DataFrame.

Example:

pandas_df = convert_to_pandas(preprocessed_df)
print(pandas_df)

Output:

   id   name                               Id
0   1   John  0a539f3c... (UUID)
1   2   Jane  1d2e4f5a... (UUID)

6. split_and_convert(vector_string)

Description: This function splits a colon-separated string and converts each part to an integer, handling edge cases.

Parameters:

  • vector_string (str): The colon-separated string to split and convert.

Returns:

  • list: The list of integers.

Example:

result = split_and_convert("1:2:3")
print(result)

Output:

[1, 2, 3]

7. add_lists(row)

Description: This function adds corresponding elements of lists contained in a row, handling different lengths.

Parameters:

  • row (list of lists): The lists to add.

Returns:

  • list: The resulting list after addition.

Example:

result = add_lists([[1, 2, 3], [4, 5, 6]])
print(result)

Output:

[5, 7, 9]

8. increment_index(row)

Description: This function creates a list of incremental integers starting from 1, with the same length as the input list.

Parameters:

  • row (list): The input list.

Returns:

  • list: The list of incremental integers.

Example:

result = increment_index([10, 20, 30])
print(result)

Output:

[1, 2, 3]

9. index(row)

Description: This function creates a list of incremental integers starting from 0, with the same length as the input list.

Parameters:

  • row (list): The input list.

Returns:

  • list: The list of incremental integers.

Example:

result = index([10, 20, 30])
print(result)

Output:

[0, 1, 2]

10. multiply_lists(row1, row2)

Description: This function multiplies corresponding elements of two lists and doubles the result.

Parameters:

  • row1 (list): The first list of numbers.
  • row2 (list): The second list of numbers.

Returns:

  • list: The list of multiplied and doubled results.

Example:

result = multiply_lists([1, 2, 3], [4, 5, 6])
print(result)

Output:

[8, 20, 36]

11. multiplyP_lists(row1, row2)

Description: This function multiplies corresponding elements of two lists.

Parameters:

  • row1 (list): The first list of numbers.
  • row2 (list): The second list of numbers.

Returns:

  • list: The list of multiplied results.

Example:

result = multiplyP_lists([1, 2, 3], [4, 5, 6])
print(result)

Output:

[4, 10, 18]

12. sum_list(row)

Description: This function calculates the sum of elements in a list.

Parameters:

  • row (list): The input list.

Returns:

  • int: The sum of the list elements.

Example:

result = sum_list([1, 2, 3])
print(result)

Output:

6

13. list_to_colon_separated_string(lst)

Description: This function converts a list of numbers to a colon-separated string.

Parameters:

  • lst (list): The list of numbers.

Returns:

  • str: The colon-separated string.

Example:

result = list_to_colon_separated_string([1, 2, 3])
print(result)

Output:

"1:2:3"

14. calculate_availability(df, file_count_col, downtime_col, period_hours, total_periods, comments_threshold)

Description: This function calculates availability based on file count and downtime, and adds availability comments.

Parameters:

  • df (DataFrame): The Spark DataFrame to process.
  • file_count_col (str): The column name for file count.
  • downtime_col (str): The column name for downtime.
  • period_hours (int): The number of hours in a period.
  • total_periods (int): The total number of periods.
  • comments_threshold (int): The threshold for adding comments.

Returns:

  • DataFrame: The DataFrame with calculated availability and comments.

Example:

result_df = calculate_availability(df, 'file_count', 'downtime', 24, 7, 10)
result_df.show()

Output:

+------------+---------+-------------+---------------------+
| file_count | downtime| Availability|Availability_Comments|
+------------+---------+-------------+---------------------+
| 10         | 5       | 99.2        |                     |
| 3          | 15      | 99.0        |Suspected Down Site  |
+------------+---------+-------------+----------------------

Project details


Download files

Download the file for your platform. If you're not sure which to choose, learn more about installing packages.

Source Distribution

pyspark-json-loader-0.1.7.tar.gz (4.7 kB view details)

Uploaded Source

Built Distribution

If you're not sure about the file name format, learn more about wheel file names.

pyspark_json_loader-0.1.7-py3-none-any.whl (5.2 kB view details)

Uploaded Python 3

File details

Details for the file pyspark-json-loader-0.1.7.tar.gz.

File metadata

  • Download URL: pyspark-json-loader-0.1.7.tar.gz
  • Upload date:
  • Size: 4.7 kB
  • Tags: Source
  • Uploaded using Trusted Publishing? No
  • Uploaded via: twine/5.1.0 CPython/3.10.12

File hashes

Hashes for pyspark-json-loader-0.1.7.tar.gz
Algorithm Hash digest
SHA256 ff4bd3ed44ccd35339a696307a5d42e880d1d9a1918a7f9bdecef474886ff6c0
MD5 b13c6f84151e29cfd66b744e09a9e021
BLAKE2b-256 8e643ab427583e3d940b87ac0a4c834206d21925316372c2f960b164a4433a16

See more details on using hashes here.

File details

Details for the file pyspark_json_loader-0.1.7-py3-none-any.whl.

File metadata

File hashes

Hashes for pyspark_json_loader-0.1.7-py3-none-any.whl
Algorithm Hash digest
SHA256 2409a21f72e3fc4b1326230eb9d7d969efd9b2e3f930cbc9818bfbcd06beec6c
MD5 b29f3d28c6abb9461b87a163246b6659
BLAKE2b-256 ad130441330fff880232137dc5ce03cb99f3fc373586a9bce6479a4c85d9346f

See more details on using hashes here.

Supported by

AWS Cloud computing and Security Sponsor Datadog Monitoring Depot Continuous Integration Fastly CDN Google Download Analytics Pingdom Monitoring Sentry Error logging StatusPage Status page