Pyspark error Invalid call to qualifier on unresolved object

Viewed 17

Hi I have table with a nested structure and trying to insert the data into the table using spark

CREATE TABLE IF NOT EXISTS dbname.tbname (
    client_id BIGINT  COMMENT 'Client id ',
    visit_id string  COMMENT 'total time spent per visit in seconds',
    visit_start_time timestamp COMMENT 'minimum timestamp per visit level',
    visit_end_time timestamp COMMENT 'maximum timestamp per visit level',
    time_spent_per_visit_in_seconds integer COMMENT 'total time spent per visit in seconds',
    visit_metrics STRUCT <
        distinct_processes_count:integer COMMENT 'Total number of distinct process browsed through the visit through web activity',
        distinct_page_count:integer COMMENT 'Total number of distinct page browsed through the visit through web activity'
        > COMMENT 'Visit Level Metrics' ,
    open_an_account STRUCT < 
        existing_client:STRUCT < 
                process_initiated:integer COMMENT 'Open an account of exisisting client intiated count through web activity',
                process_submitted:integer COMMENT 'Open an account of exisisting client submitted count through web activity'
                > COMMENT 'Open and account counts through web activity' ,
        unspecified:STRUCT < 
                process_initiated:integer COMMENT 'Open an account of client intiated initiated count through web activity',
                process_submitted:integer COMMENT 'Open an account of client intiated initiated count through web activity'
                > COMMENT 'Open and account counts through web activity' 
        > COMMENT 'Open and account web activity count metrics',
    transfer_of_assets STRUCT < 
        process_initiated:integer COMMENT 'Transfer of assests start indicator through web activity',
        process_submitted:integer COMMENT 'Transfer of assests end indicator through web activity'
        > COMMENT 'Transfer of account metrics',
    direct_rollover STRUCT < 
        process_initiated:integer COMMENT 'Direct Rollover process intitiation count through web activity',
        process_submitted:integer COMMENT 'Direct Rollover process submission count through web activity'
        > COMMENT 'Direct Rollover Metrics',
    trade_path STRUCT <
        process_initiated:integer COMMENT 'Trade Path process intitiation count through web activity',
        process_submitted:integer COMMENT 'Trade Path process intitiation count through web activity' 
        > COMMENT 'Trade Path Process Metrics',
    account_management STRUCT < 
        account_activity:STRUCT < 
            alert_change_counts:STRUCT < 
                process_initiated:integer COMMENT 'Account Management AccountActivityAlerts change edit initiated count  through web activity',
                process_submitted:integer COMMENT 'Account Management AccountActivityAlerts change edit submitted  through web activity'
                > COMMENT 'Account Activity Alert Change counts through web activity' 
        > COMMENT 'Account Management Account Activity' ,
        contact_info:STRUCT < 
            address_or_phone_change:STRUCT < 
                process_initiated:integer COMMENT 'Account Management addressphone change initiated count  through web activity',
                process_submitted:integer COMMENT 'Account Management addressphone change submitted count  through web activity'
                > COMMENT 'Account Activity Alert Change counts',
            email_change:STRUCT < 
                process_initiated:integer COMMENT 'Account Management email change initiated count through web activity',
                process_submitted:integer COMMENT 'Account Management email change submitted count through web activity'
                > COMMENT 'Account Activity Email activity counts', 
            username_password_change:STRUCT < 
                process_initiated:integer COMMENT 'Account Management UsernamePassword change initiated count through web activity',
                process_submitted:integer COMMENT 'Account Management UsernamePassword change submitted count through web activity'
                > COMMENT 'Account Activity Username Password activity counts'       
        > COMMENT 'Account Management Contact Info Changes' ,
        automatic_investments:STRUCT < 
            add_operation:STRUCT < 
                process_initiated:integer COMMENT 'Account Management Auto investment Add initiated count through web activity',
                process_submitted:integer COMMENT 'Account Management Auto investment Add submitted count through web activity'
                > COMMENT 'Account Activity add Auto investment counts',
            edit_operation:STRUCT < 
                process_initiated:integer COMMENT 'Account Management Auto investment initiated count through web activity',
                process_submitted:integer COMMENT 'Account Management Auto investment submitted count through web activity'
                > COMMENT 'Account Activity edit Auto investment counts',   
            delete_operation:STRUCT < 
                process_initiated:integer COMMENT 'Account Management Auto investment delete initiated count through web activity',
                process_submitted:integer COMMENT 'Account Management Auto investment delete submitted count through web activity'
                > COMMENT 'Account Management delete Auto investment counts'         
        > COMMENT 'Account Management Auto investment Changes' ,
        automatic_withdrawals:STRUCT < 
                process_initiated:integer COMMENT 'Account Management Auto withdrawal initiated count',
                process_submitted:integer COMMENT 'Account Management Auto withdrawal submitted count'
        > COMMENT 'Account Management Auto withdrawal process counts' ,
        banking_info_addition:STRUCT < 
                process_initiated:integer COMMENT 'Account Management bank information initiated count',
                process_submitted:integer COMMENT 'Account Management bank information submitted count'
        > COMMENT 'Account Management bank information addition process counts' ,
        beneficiaries_change:STRUCT < 
                process_initiated:integer COMMENT 'Account Management beneficiaries change initiated count',
                process_submitted:integer COMMENT 'Account Management beneficiaries change submitted count'
        > COMMENT 'Account Management beneficiaries change counts' ,
        cost_basis_method_change:STRUCT < 
            process_initiated:integer COMMENT 'Account Management Costbasis method change initiated count',
            process_submitted:integer COMMENT 'Account Management Costbasis method change submitted count'
        > COMMENT 'Account Management Costbasis method counts' ,
        holding_dividends_capital_gains_change:STRUCT < 
            process_initiated:integer COMMENT 'Account Management HoldingDividendsCapitalGains initiated count',
            process_submitted:integer COMMENT 'Account Management HoldingDividendsCapitalGains submitted count'
        > COMMENT 'Account Management HoldingDividendsCapitalGains counts' ,
        instantbank_change:STRUCT < 
            process_initiated:integer COMMENT 'Account Management  instant bank add  initiated count',
            process_submitted:integer COMMENT 'Account Management  instant bank add  submitted count'
        > COMMENT 'Account Management HoldingDividendsCapitalGains counts' ,
        required_minimum_distribution:STRUCT < 
            calculate:STRUCT < 
                process_initiated:integer COMMENT 'Account Management required minimum distribution calculate initiated count through web activity',
                process_submitted:integer COMMENT 'Account Management required minimum distribution calculate submitted count through web activity'
                > COMMENT 'Account Activity Alert Change counts',
            one_time_distribution:STRUCT < 
                process_initiated:integer COMMENT 'Account Management  required minimum distribution OneTimeDistribution intiated through web activity',
                process_submitted:integer COMMENT 'Account Management  required minimum distribution OneTimeDistribution submitted through web activity'
                > COMMENT 'Account Activity Email activity counts' 
        > COMMENT 'Account Management Activity Changes'
    > COMMENT ' Account Management Metrics'
    )
    COMMENT 'All journey visits provides a different Journey process visited during each visit an client made'
    PARTITIONED BY (visit_date date COMMENT 'Date of the visit. Partition column')
    STORED AS PARQUET
    LOCATION 's3://bucketname/db_name/table_name/'

The loading of the data logic is as below i am just restructuring using the named_struct function nothing much than that.

I has to use sql to do this as the frameworkwe have accepts only sql and not dataframes.

spark.sql('''with web_cust_copy as 
(
   select
    cast ('1' as bigint) client_id,
    '213213' visit_id,
    cast ('10' as integer) as time_spent_per_visit_in_seconds,
    '2022-09-17 00:01:23' as visit_start_time,
    '2022-09-17 00:01:23'  as visit_end_time,
    cast ('0' as integer) total_distinct_page_count ,
    cast ('2' as integer) total_distinct_process_count,
    cast ('1' as integer ) OpenanAccountJourney_initiated,
    cast ('1' as integer ) OpenanAccountJourney_submitted,
    cast ('1' as integer ) Client_OpenanAccountJourney_initiated,
    cast ('1'as integer ) Client_OpenanAccountJourney_submitted,
    cast ('1' as integer ) TransferofAssets_initiated,
    cast ('1' as integer ) TransferofAssets_submitted,
    cast ('1' as integer ) DirectRollover_initiated,
    cast ('1' as integer ) DirectRollover_submitted,
    cast ('1' as integer ) Trade_Path_initiated,
    cast ('1' as integer ) Trade_Path_submitted,
    cast ('1' as integer ) AM_UserNamePassword_Change_initiated,
    cast ('1' as integer ) AM_UserNamePassword_Change_submitted,
    cast ('1' as integer ) AM_RMD_OneTimeDistribution_initiated,
    cast ('1' as integer ) AM_RMD_OneTimeDistribution_submitted,
    cast ('1' as integer ) AM_RMD_Calculate_initiated,
    cast ('1' as integer ) AM_RMD_Calculate_submitted,
    cast ('1' as integer ) AM_InstantBank_Add_initiated,
    cast ('1' as integer ) AM_InstantBank_Add_submitted,
    cast ('1' as integer ) AM_HoldingDividendsCapitalGains_Change_initiated,
    cast ('1' as integer ) AM_HoldingDividendsCapitalGains_Change_submitted,
    cast ('1' as integer ) AM_Email_Change_initiated,
    cast ('1' as integer ) AM_Email_Change_submitted,
    cast ('1' as integer ) AM_CostBasisMethod_Change_initiated,
    cast ('1' as integer ) AM_CostBasisMethod_Change_submitted,
    cast ('1' as integer ) AM_Beneficiaries_Change_initiated,
    cast ('1' as integer ) AM_Beneficiaries_Change_submitted,
    cast ('1' as integer ) AM_BankInformation_Add_initiated,
    cast ('1' as integer ) AM_BankInformation_Add_submitted,
    cast ('1' as integer ) AM_AutomaticWithdrawal_Add_initiated,
    cast ('1' as integer ) AM_AutomaticWithdrawal_Add_submitted,
    cast ('1' as integer ) AM_AutomaticInvestment_Edit_initiated,
    cast ('1' as integer ) AM_AutomaticInvestment_Edit_submitted,
    cast ('1' as integer ) AM_AutomaticInvestment_Delete_intitated,
    cast ('1' as integer ) AM_AutomaticInvestment_Delete_submitted,
    cast ('1' as integer ) AM_AutomaticInvestment_Add_intitated,
    cast ('1' as integer ) AM_AutomaticInvestment_Add_submitted,
    cast ('1' as integer ) AM_AddressPhone_Change_intitated,
    cast ('1' as integer ) AM_AddressPhone_Change_submitted,
    cast ('1' as integer ) AM_AccountActivityAlerts_Change_intitated,
    cast ('1' as integer ) AM_AccountActivityAlerts_Change_submitted,
    cast ('2022-09-17' as date) visit_date
 ),

all_journey_visits_data_struct as 

(

SELECT client_id
,visit_id
,visit_start_time
,visit_end_time
,time_spent_per_visit_in_seconds
,named_struct('distinct_processes_count',total_distinct_process_count,'distinct_page_count',total_distinct_page_count) as visit_metrics
,named_struct('process_initiated',openanaccountjourney_initiated,'process_submitted',openanaccountjourney_submitted) as open_an_acct_journey
,named_struct('process_initiated',client_openanaccountjourney_initiated,'process_submitted',client_openanaccountjourney_submitted) as client_open_an_acct_journey
,named_struct('process_initiated',transferofassets_initiated,'process_submitted',transferofassets_submitted) as transfer_of_assets
,named_struct('process_initiated',directrollover_initiated,'process_submitted',directrollover_submitted) as direct_rollover
,named_struct('process_initiated',trade_path_initiated,'process_submitted',trade_path_submitted) as trade_path
,named_struct('process_initiated',am_accountactivityalerts_change_intitated,'process_submitted',am_accountactivityalerts_change_submitted) as alert_change_counts
,named_struct('process_initiated',am_email_change_initiated,'process_submitted',am_email_change_submitted) as email_change
,named_struct('process_initiated',am_usernamepassword_change_initiated,'process_submitted',am_usernamepassword_change_submitted) as username_password_change
,named_struct('process_initiated',am_addressphone_change_intitated,'process_submitted',am_addressphone_change_submitted) as address_or_phone_change
,named_struct('process_initiated',am_automaticinvestment_add_intitated,'process_submitted',am_automaticinvestment_add_submitted) as automatic_investments_add_operation
,named_struct('process_initiated',am_automaticinvestment_edit_initiated,'process_submitted',am_automaticinvestment_edit_submitted) as automatic_investments_edit_operation
,named_struct('process_initiated',am_automaticinvestment_delete_intitated,'process_submitted',am_automaticinvestment_delete_submitted) as automatic_investments_delete_operation
,named_struct('process_initiated',am_automaticwithdrawal_add_initiated,'process_submitted',am_automaticwithdrawal_add_submitted) as automatic_withdrawals
,named_struct('process_initiated',am_bankinformation_add_initiated,'process_submitted',am_bankinformation_add_submitted) as banking_info_addition
,named_struct('process_initiated',am_beneficiaries_change_initiated,'process_submitted',am_beneficiaries_change_submitted) as beneficiaries_change
,named_struct('process_initiated',am_costbasismethod_change_initiated,'process_submitted',am_costbasismethod_change_submitted) as cost_basis_method_change
,named_struct('process_initiated',am_holdingdividendscapitalgains_change_initiated,'process_submitted',am_holdingdividendscapitalgains_change_submitted) as holding_dividends_capital_gains_change
,named_struct('process_initiated',am_instantbank_add_initiated,'process_submitted',am_instantbank_add_submitted) as instantbank_change
,named_struct('process_initiated',am_rmd_calculate_initiated,'process_submitted',am_rmd_calculate_submitted) as rmd_calculate
,named_struct('process_initiated',am_rmd_onetimedistribution_initiated,'process_submitted',am_rmd_onetimedistribution_submitted) as rmd_one_time_distribution
,visit_date from web_cust_copy where visit_date = '2022-09-17'
)
INSERT OVERWRITE TABLE dbname.tbname partition(visit_date)
SELECT /*+ COALESCE(4) */ client_id
,visit_id
,visit_start_time
,visit_end_time
,time_spent_per_visit_in_seconds
, visit_metrics
,named_struct('existing_client',client_open_an_acct_journey,'unspecified',open_an_acct_journey) as open_an_account
,transfer_of_assets
,direct_rollover
,trade_path
,named_struct(
   'account_activity',(named_struct('alert_change_counts',alert_change_counts))
   ,'contact_info',(named_struct('address_or_phone_change',address_or_phone_change,'email_change',email_change,'username_password_change',username_password_change))
   ,'automatic_investments',(named_struct('add_operation',automatic_investments_add_operation,'edit_operation',automatic_investments_edit_operation,'delete_operation',automatic_investments_delete_operation))
   ,'automatic_withdrawals',(named_struct('automatic_withdrawals',automatic_withdrawals))
   ,'banking_info_addition',(named_struct('banking_info_addition',banking_info_addition))
   ,'beneficiaries_change',(named_struct('beneficiaries_change',beneficiaries_change))
   ,'cost_basis_method_change',(named_struct('cost_basis_method_change',beneficiaries_change))
   ,'holding_dividends_capital_gains_change',(named_struct('holding_dividends_capital_gains_change',holding_dividends_capital_gains_change))
   ,'instantbank_change',(named_struct('instantbank_change',instantbank_change))
   ,'required_minimum_distribution',(named_struct('calculate',rmd_calculate,'one_time_distribution',rmd_one_time_distribution))
   ) as account_management
,visit_date
FROM all_journey_visits_data_struct''')

Error iam encountering is below

An error was encountered:
"Invalid call to qualifier on unresolved object, tree: 'account_management"
Traceback (most recent call last):
  File "/mnt/yarn/usercache/umh7/appcache/application_1663754629436_0006/container_1663754629436_0006_01_000001/pyspark.zip/pyspark/sql/session.py", line 767, in sql
    return DataFrame(self._jsparkSession.sql(sqlQuery), self._wrapped)
  File "/mnt/yarn/usercache/appcache/application_1663754629436_0006/container_1663754629436_0006_01_000001/py4j-0.10.7-src.zip/py4j/java_gateway.py", line 1257, in __call__
    answer, self.gateway_client, self.target_id, self.name)
  File "/mnt/yarn/usercache/appcache/application_1663754629436_0006/container_1663754629436_0006_01_000001/pyspark.zip/pyspark/sql/utils.py", line 71, in deco
    raise AnalysisException(s.split(': ', 1)[1], stackTrace)
pyspark.sql.utils.AnalysisException: "Invalid call to qualifier on unresolved object, tree: 'account_management"

Not sure how to fix the error

0 Answers
Related