id: automated-customer-feedback-sentiment-analysis
namespace: company.customer_success
inputs:
- id: file_url
type: STRING
defaults: https://huggingface.co/datasets/kestra/datasets/raw/main/csv/customer_feedback.csv
description: URL of the customer feedback CSV file to analyze
tasks:
- id: download_feedback
type: io.kestra.plugin.core.http.Download
uri: "{{ inputs.file_url }}"
- id: analyze_sentiment
type: io.kestra.plugin.scripts.python.Script
warningOnStdErr: false
taskRunner:
type: io.kestra.plugin.scripts.runner.docker.Docker
containerImage: ghcr.io/kestra-io/pydata:latest
env:
OPENAI_API_KEY: "{{ secret('OPENAI_API_KEY') }}"
script: |
import pandas as pd
import json
from openai import OpenAI
import os
# Initialize OpenAI client
client = OpenAI(api_key=os.environ["OPENAI_API_KEY"])
# Read the feedback data
df = pd.read_csv("{{ outputs.download_feedback.uri }}")
# Function to analyze sentiment using OpenAI
def get_sentiment(text):
try:
response = client.chat.completions.create(
model="gpt-4o-mini",
messages=[
{"role": "system", "content": "You are a customer success AI. Classify the sentiment of the following customer feedback as POSITIVE, NEUTRAL, or NEGATIVE. Also identify the primary topic (e.g., Pricing, Support, Feature Request, Bug). Respond strictly in JSON format: {\"sentiment\": \"...\", \"topic\": \"...\"}"},
{"role": "user", "content": text}
],
response_format={"type": "json_object"}
)
return json.loads(response.choices[0].message.content)
except Exception as e:
return {"sentiment": "ERROR", "topic": "ERROR"}
# We only process the top 20 recent feedbacks to save API costs
recent_feedbacks = df.head(20).copy()
results = []
for idx, row in recent_feedbacks.iterrows():
analysis = get_sentiment(row['feedback_text'])
results.append({
"id": row.get('id', idx),
"feedback": row['feedback_text'],
"sentiment": analysis['sentiment'],
"topic": analysis['topic']
})
# Save results to a new JSON file
output_df = pd.DataFrame(results)
output_df.to_json("sentiment_results.json", orient="records")
# Calculate quick aggregates for Kestra metrics
positive_count = len(output_df[output_df['sentiment'] == 'POSITIVE'])
negative_count = len(output_df[output_df['sentiment'] == 'NEGATIVE'])
print(f"::setOutputs {{\"positive_count\": {positive_count}, \"negative_count\": {negative_count}}}")
outputFiles:
- "sentiment_results.json"
- id: aggregate_results_duckdb
type: io.kestra.plugin.jdbc.duckdb.Query
inputFiles:
data.json: "{{ outputs.analyze_sentiment.outputFiles['sentiment_results.json'] }}"
sql: |
INSTALL json;
LOAD json;
CREATE TABLE feedback AS SELECT * FROM read_json_auto('data.json');
COPY (
SELECT
topic,
COUNT(*) as total_feedback,
SUM(CASE WHEN sentiment = 'POSITIVE' THEN 1 ELSE 0 END) as positive_reviews,
SUM(CASE WHEN sentiment = 'NEGATIVE' THEN 1 ELSE 0 END) as negative_reviews
FROM feedback
GROUP BY topic
ORDER BY negative_reviews DESC
) TO '{{ outputDir }}/aggregated_report.csv' WITH (HEADER 1, DELIMITER ',');
SELECT topic, COUNT(*) as issues FROM feedback WHERE sentiment = 'NEGATIVE' GROUP BY topic;
fetchType: FETCH
- id: send_slack_alert
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"text": "📊 *Daily Customer Feedback Analysis Complete*",
"blocks": [
{
"type": "section",
"text": {
"type": "mrkdwn",
"text": "*Daily Customer Feedback Analysis Complete*\n\nWe analyzed the latest customer feedback. Here is the summary:\n• 🟢 *Positive Feedback:* {{ outputs.analyze_sentiment.vars.positive_count }}\n• 🔴 *Negative Feedback:* {{ outputs.analyze_sentiment.vars.negative_count }}\n\nThe most pressing negative topics are:\n{% for row in outputs.aggregate_results_duckdb.rows %}\n> *{{ row.topic }}*: {{ row.issues }} negative mentions\n{% endfor %}"
}
}
]
}