please can anybody help i have been stuck for 4 da...
# random
r
please can anybody help i have been stuck for 4 days now!!!!!!!!!!alert the code i used was properly working fine with a namespace that didnt have authentication setup. when i did the authorization only i am facing this issue. the schema creation and table creation is working fine after authentication done in pinot. but the data ingestion is having issues. these are the config files: when i create the schema and table for a sample table. when i try ingest a sample csv file i cant do that: i am getting the above mentioned error in the chat. Upload response: 500 Error: {"code":500,"error":"Caught exception when ingesting file into table: test_tenant_OFFLINE. Caught exception while uploading segments. Push mode: TAR, segment tars: [[file:/opt/pinot/tmp/ingestion_dir/working_dir_test_tenant_OFFLINE_1742312090194_FRKYyJJrdp/segment_tar_dir/test_tenant_1742312090217.tar.gz]]"} ingest-pythoncode:
Copy code
import requests
import json
from pathlib import Path
import time

PINOT_CONTROLLER = "********"
PINOT_BROKER = "*********"
AUTH_HEADERS = {
    "Authorization": "Basic YWRtaW46dmVyeXNlY3JldA==",
    "Content-Type": "application/json"
}

def verify_segment(tenant_name, segment_name):
    max_retries = 10
    retry_interval = 2  # seconds
    
    for i in range(max_retries):
        print(f"\nChecking segment status (attempt {i+1}/{max_retries})...")
        
        response = requests.get(f"{PINOT_CONTROLLER}/segments/{tenant_name}/{segment_name}/metadata",headers=AUTH_HEADERS)
        if response.status_code == 200:
            print("Segment is ready!")
            return True
            
        print(f"Segment not ready yet, waiting {retry_interval} seconds...")
        time.sleep(retry_interval)
    
    return False

def simple_ingest(tenant_name):
    # Verify schema and table exist
    schema_response = requests.get(f'{PINOT_CONTROLLER}/schemas/{tenant_name}',headers=AUTH_HEADERS)
    table_response = requests.get(f'{PINOT_CONTROLLER}/tables/{tenant_name}',headers=AUTH_HEADERS)
    
    if schema_response.status_code != 200 or table_response.status_code != 200:
        print(f"Schema or table missing for tenant {tenant_name}. Please run create_schema.py and create_table.py first")
        return

    csv_path = Path(f"data/{tenant_name}_data.csv")
    
    print(f"\nUploading data for tenant {tenant_name}...")
    with open(csv_path, 'rb') as f:
        files = {'file': (f'{tenant_name}_data.csv', f, 'text/csv')}
        
        # Using a dictionary for column mapping first
        column_map = {
            "Ticket ID": "Ticket_ID",
            "Customer Name": "Customer_Name",
            "Customer Email": "Customer_Email",
            "Company_name": "Company_name",
            "Customer Age": "Customer_Age",
            "Customer Gender": "Customer_Gender",
            "Product purchased": "product_purchased",
            "Date of Purchase": "Date_of_Purchase",
            "Ticket Subject": "Ticket_Subject",
            "Description": "Description",
            "Ticket Status": "Ticket_Status",
            "Resolution": "Resolution",
            "Ticket Priority": "Ticket_Priority",
            "Source": "Source",
            "Created date": "Created_date",
            "First Response Time": "First_Response_Time",
            "Time to Resolution": "Time_of_Resolution",
            "Number of conversations": "Number_of_conversations",
            "Customer Satisfaction Rating": "Customer_Satisfaction_Rating",
            "Category": "Category",
            "Intent": "Intent",
            "Type": "Type",
            "Relevance": "Relevance",
            "Escalate": "Escalate",
            "Sentiment": "Sentiment",
            "Tags": "Tags",
            "Agent": "Agent",
            "Agent politeness": "Agent_politeness",
            "Agent communication": "Agent_communication",
            "Agent Patience": "Agent_Patience",
            "Agent overall performance": "Agent_overall_performance",
            "Conversations": "Conversations"
        }

        config = {
            "inputFormat": "csv",
            "header": "true",
            "delimiter": ",",
            "fileFormat": "csv",
            "multiValueDelimiter": ";",
            "skipHeader": "false",
            # Convert dictionary to proper JSON string
            "columnHeaderMap": json.dumps(column_map)
        }
        
        params = {
            'tableNameWithType': f'{tenant_name}_OFFLINE',
            'batchConfigMapStr': json.dumps(config)
        }

        upload_headers = {
            "Authorization": AUTH_HEADERS["Authorization"]
        }
        
        response = <http://requests.post|requests.post>(
            f'{PINOT_CONTROLLER}/ingestFromFile',
            files=files,
            params=params,
            headers=upload_headers
        )
    
    print(f"Upload response: {response.status_code}")
    if response.status_code != 200:
        print(f"Error: {response.text}")
        return

    if response.status_code == 200:
        try:
            response_data = json.loads(response.text)
            print(f"Response data: {response_data}")
            segment_name = response_data["status"].split("segment: ")[1]
            print(f"\nWaiting for segment {segment_name} to be ready...")
            
            if verify_segment(tenant_name, segment_name):
                query = {
                    "sql": f"SELECT COUNT(*) FROM {tenant_name}_OFFLINE",
                    "trace": False
                }
                
                query_response = <http://requests.post|requests.post>(
                    f"{PINOT_BROKER}/query/sql",
                    json=query,
                    headers=AUTH_HEADERS
                )
                
                print("\nQuery response:", query_response.status_code)
                if query_response.status_code == 200:
                    print(json.dumps(query_response.json(), indent=2))
            else:
                print("Segment verification timed out")
        except Exception as e:
            print(f"Error processing segment: {e}")
            print(f"Full response text: {response.text}")

if __name__ == "__main__":
    simple_ingest("test_tenant")
Slack Conversation