PySpark新手求助:如何将JSON RDD解析为指定结构的DataFrame
Hey there! Let's get your PySpark DataFrame sorted out. The issue here is that the CoinMarketCap API returns a nested JSON structure, so we need to explicitly pull out the date_added and price fields instead of just reading the raw JSON directly. Here are two straightforward approaches to fix this:
Approach 1: Preprocess the JSON Data Before Creating RDD
This method extracts the required fields first using Python's JSON parsing, then creates a DataFrame from the simplified data. It's great if you want to reduce the data size passed to Spark upfront:
from pyspark import SparkConf, SparkContext from pyspark.sql import SQLContext import requests from requests.exceptions import ConnectionError, Timeout, TooManyRedirects # Initialize Spark context conf = SparkConf().setAppName('rates').setMaster("local") sc = SparkContext(conf=conf) sqlContext = SQLContext(sc) # API details url = 'https://pro-api.coinmarketcap.com/v1/cryptocurrency/quotes/latest' parameters = {'symbol': 'BTC', 'convert':'JPY'} headers = { 'Accepts': 'application/json', 'X-CMC_PRO_API_KEY': '***********************', # Replace with your actual API key } try: # Fetch and parse the API response response = requests.get(url, params=parameters, headers=headers) raw_data = response.json() # Extract only the fields we need from the nested structure btc_entry = raw_data['data']['BTC'] cleaned_data = [{ 'date_added': btc_entry['date_added'], 'price': btc_entry['quote']['JPY']['price'] }] # Create RDD and convert to DataFrame rdd = sc.parallelize(cleaned_data) target_df = sqlContext.read.json(rdd) # Optional: Convert date_added to proper date type from pyspark.sql.functions import to_date, col target_df = target_df.withColumn("date_added", to_date(col("date_added"))) target_df.show() except (ConnectionError, Timeout, TooManyRedirects) as e: print(f"API request failed: {e}") finally: # Always stop the Spark context when done sc.stop()
Approach 2: Extract Fields from the Raw Spark DataFrame
If you already have the raw json_df from your original code, you can use Spark's column functions to dig into the nested structures and rename the fields:
from pyspark.sql.functions import col, to_date # Assuming json_df is the DataFrame from your original code target_df = json_df.select( col("data.BTC.date_added").alias("date_added"), col("data.BTC.quote.JPY.price").alias("price") ) # Optional: Format the date string into a proper date type target_df = target_df.withColumn("date_added", to_date(col("date_added"))) target_df.show()
Why Your Original Code Didn't Work
The raw JSON response has a top-level structure with data and status fields. Inside data, there's a BTC struct that contains date_added, and a nested quote.JPY struct with the price. Spark reads this as nested columns by default, so we need to explicitly reference these nested paths to get the flat, clean fields you want.
Either approach will give you the DataFrame you're targeting, with just date_added and price columns.
内容的提问来源于stack exchange,提问作者delalma

