You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

PySpark新手求助:如何将JSON RDD解析为指定结构的DataFrame

How to Extract Nested Fields from CoinMarketCap API Response in PySpark

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.11 07:42:18