Weiter zum Inhalt

Datenpipelines mit R aufbauen

Lerne, wie du mit R und SQLite von Grund auf eine ETL-Pipeline erstellst, Tweets in Echtzeit sammelst und für spätere Analysen speicherst.
Aktualisiert 18. Sept. 2026  · 13 Min. lesen

Mit KI erkunden

ChatGPTClaudePerplexity

Man kann durchaus sagen: Saubere ETL-Pipelines sind das Herzstück der Data Science. Ohne aufbereitete, strukturierte Daten ist es schwer, belastbare Insights zu liefern, die bessere Geschäftsentscheidungen ermöglichen.

In diesem Tutorial bauen wir deshalb eine einfache ETL-Pipeline, die Tweets in Echtzeit direkt in eine SQLite-Datenbank mit R streamt. Das ist zum Beispiel eine gängige Aufgabe in der Analyse sozialer Netzwerke.

Im Fokus stehen der Denkprozess rund um Datenerfassung und -speicherung sowie der Umgang mit der Twitter API über das R-Paket rtweet.

Los geht’s mit den richtigen Tools. Zuerst brauchst du Zugriff auf die Twitter API. Insgesamt folgst du diesen Schritten:

  • Erstelle einen Twitter-Account, falls noch nicht vorhanden.
  • Folge diesem Link und beantrage einen Developer-Account (Hinweis: Twitter muss die Bewerbung mittlerweile genehmigen).
  • Erstelle eine neue App auf dieser Webseite.
  • Trage alle App-Details ein und generiere deinen Access Token.
  • Sichere dir Consumer Key, Consumer Secret, Access Token und Access Token Secret – damit stellst du die Verbindung zur API her.

Wenn der Twitter-Zugang steht, brauchst du noch SQLite, sofern nicht vorhanden. Eine vollständige Anleitung zur Installation findest du im Beginner’s Guide to SQLite hier bei DataCamp. Für dieses Tutorial nutzen wir SQLite, weil es besonders einfach zu handhaben ist.

Schritt 1: Datenbank und Tabelle zum Speichern der Twitter-Daten erstellen

Mit Zugriff auf die Twitter API und installierter SQLite können wir endlich die Pipeline bauen, um Tweets beim Streamen dauerhaft abzulegen. Zunächst legen wir die neue SQLite-Datenbank in R an:

# Import necessary libraries and functions
library(RSQLite)
library(rtweet)
library(tm)
library(dplyr)
library(knitr)
library(wordcloud)
library(lubridate)
library(ggplot2)
source("transform_and_clean_tweets.R")
# Create our SQLite database
conn <- dbConnect(RSQLite::SQLite(), "Tweet_DB.db")

Als Nächstes erstellen wir in der Datenbank eine Tabelle für die Tweets. In unserem Beispiel speichern wir folgende Variablen:

  • Tweet_ID als INTEGER Primary Key
  • User als TEXT
  • Tweet_Content als TEXT
  • Date_Created als INTEGER

Warum speichern wir Datumswerte als Integer? SQLite hat keinen eigenen Datentyp für Datum und Uhrzeit. Deshalb speichern wir die Anzahl Sekunden seit dem 01.01.1970.

Jetzt schreiben wir die Tabelle:

dbExecute(conn, "CREATE TABLE Tweet_Data(
                  Tweet_ID INTEGER PRIMARY KEY,
                  User TEXT,
                  Tweet_Content TEXT,
                  Date_Created INTEGER)")

Nachdem die Tabelle erstellt ist, kannst du in sqlite3.exe prüfen, ob sie vorhanden ist. So sieht das zum Beispiel aus:

sqlite3.exe Screenshot


Schritt 2: Tweets zu deinen Lieblingsthemen streamen

Glaub es oder nicht: Damit sind die Anforderungen und die nötige Infrastruktur für eine einfache, funktionierende Twitter-Streaming-Pipeline bereits erfüllt. Jetzt streamen wir Tweets über die API. Hinweis: Für dieses Tutorial nutze ich die kostenlose Standard-API. Es gibt kostenpflichtige Premium-Varianten, die für Forschung oder höhere Anforderungen besser passen können.

Ohne weitere Umschweife richten wir unseren Twitter-Listener ein. Zuerst importierst du das rtweet-Paket und trägst die Access Tokens und Secrets deiner App wie eingangs beschrieben ein:

token <- create_token(app = 'Your_App_Name',
                      consumer_key = 'Your_Consumer_Key',
                      consumer_secret = 'Your_Consumer_Secret',
                      access_token = 'Your_Access_Token',
                      access_secret = 'Your_Access_Secret')

Mit dem Token entscheidest du als Nächstes, welche Tweets du streamen willst. Die Funktion stream_tweets im rtweet-Paket bietet viele Optionen für Abfragen an die Twitter API. Du kannst zum Beispiel Tweets streamen, die bestimmte Hashtags oder Keywords enthalten (bis zu 400), einen kleinen Zufallsausschnitt aller öffentlichen Tweets, die Tweets bestimmter User-IDs oder Screennamen (bis zu 5000) verfolgen oder Tweets über Geokoordinaten einsammeln.

Für dieses Tutorial streame ich Tweets mit Hashtags rund um Data Science (siehe Liste unten). Das Format, in dem die Hashtags angegeben sind, wirkt etwas ungewohnt, ist aber genau so von stream_tweets gefordert, wenn du nach Hashtags oder Keywords horchst. Für Nutzerlisten oder Koordinaten weicht das Format ab. Details findest du in der Dokumentation.

keys <- "#nlp,#machinelearning,#datascience,#chatbots,#naturallanguageprocessing,#deeplearning"

Mit den Keywords definiert, bauen wir nun die Streaming-Schleife. Es gibt mehrere Ansätze; dieses Muster hat sich für mich bewährt:

# Initialize the streaming hour tally
hour_counter <- 0

# Initialize a while loop that stops when the number of hours you want to stream tweets for is exceeded
while(hour_counter <= 12){
  # Set the stream time to 2 hours each iteration (7200 seconds)
  streamtime <- 7200
  # Create the file name where the 2 hour stream will be stored. Note that the Twitter API outputs a .json file.
  filename <- paste0("nlp_stream_",format(Sys.time(),'%d_%m_%Y__%H_%M_%S'),".json")
  # Stream Tweets containing the desired keys for the specified amount of time
  stream_tweets(q = keys, timeout = streamtime, file_name = filename)
  # Clean the streamed tweets and select the desired fields
  clean_stream <- transform_and_clean_tweets(filename, remove_rts = TRUE)
  # Append the streamed tweets to the Tweet_Data table in the SQLite database
  dbWriteTable(conn, "Tweet_Data", clean_stream, append = T)
  # Delete the .json file from this 2-hour stream
  file.remove(filename)
  # Add the hours to the tally
  hour_counter <- hour_counter + 2
}

Diese Schleife streamt in 2‑Stunden-Intervallen insgesamt 12 Stunden lang so viele Tweets wie möglich, die einen der Hashtags im String enthalten. Alle 2 Stunden legt der Listener im aktuellen Arbeitsverzeichnis eine .json-Datei mit dem in filename definierten Namen ab.

Anschließend übergibt sie den Dateinamen an die Funktion transform_and_clean_tweets, die auf Wunsch Retweets entfernt, die gewünschten Spalten aus der API-Antwort auswählt und den Tweet-Text normalisiert.

Danach hängt sie das resultierende Dataframe an die Tweet_Data-Tabelle in unserer SQLite-Datenbank an. Zum Schluss wird der Stundenzähler um 2 erhöht (weil jeder Stream 2 Stunden dauert) und die erzeugte .json-Datei gelöscht. Das ist sinnvoll, weil alle relevanten Daten inzwischen in der Datenbank liegen und die .json-Dateien sonst Speicher fressen.

Schauen wir uns die Funktion transform_and_clean_tweets im Detail an:

transform_and_clean_tweets <- function(filename, remove_rts = TRUE){

  # Import the normalize_text function
  source("normalize_text.R")

  # Parse the .json file given by the Twitter API into an R data frame
  df <- parse_stream(filename)
  # If remove_rst = TRUE, filter out all the retweets from the stream
  if(remove_rts == TRUE){
    df <- filter(df,df$is_retweet == FALSE)
  }
  # Keep only the tweets that are in English
  df <- filter(df, df$lang == "en")
  # Select the features that you want to keep from the Twitter stream and rename them
  # so the names match those of the columns in the Tweet_Data table in our database
  small_df <- df[,c("screen_name","text","created_at")]
  names(small_df) <- c("User","Tweet_Content","Date_Created")
  # Finally normalize the tweet text
  small_df$Tweet_Content <- sapply(small_df$Tweet_Content, normalize_text)
  # Return the processed data frame
  return(small_df)
}

Wie beschrieben filtert die Funktion bei Bedarf Retweets heraus, behält die relevanten Felder und normalisiert den Text. Im Grunde bildet sie das „T“ in ETL (Transformation). Ein zentrales Element ist dabei das Säubern des Tweet-Texts.

Textdaten brauchen meist einige Vorverarbeitungsschritte, bevor sie analysierbar sind. Bei Tweets kann das das Entfernen von URLs, Stoppwörtern und Erwähnungen, Kleinschreibung, Stemming usw. umfassen. Nicht alles ist immer nötig. Hier ist zunächst die von mir verwendete normalize_text-Funktion zur Vorverarbeitung:

normalize_text <- function(text){
  # Keep only ASCII characters
  text = iconv(text, "latin1", "ASCII", sub="")
  # Convert to lower case characters
  text = tolower(text)
  # Remove any HTML tags
  text = gsub("<.*?>", " ", text)
  # Remove URLs
  text = gsub("\\s?(f|ht)(tp)(s?)(://)([^\\.]*)[\\.|/](\\S*)", "", text)
  # Keep letters and numbers only
  text = gsub("[^[:alnum:]]", " ", text)
  # Remove stop words
  text = removeWords(text,c("rt","gt",stopwords("en")))
  # Remove any extra white space
  text = stripWhitespace(text)                                 
  text = gsub("^\\s+|\\s+$", "", text)                         

  return(text)
}

Je nach Use Case können diese Schritte reichen. Wie erwähnt, kannst du weitere Schritte wie Stemming oder Lemmatisierung ergänzen oder nur Buchstaben statt Buchstaben und Zahlen zulassen. Experimentiere ruhig mit verschiedenen Kombinationen – ideal, um Regex-Fähigkeiten zu trainieren.

Nach all diesen Schritten liegt eine SQLite-Datenbank vor, die alle gestreamten Tweets enthält. Ob alles funktioniert hat, prüfst du mit ein paar einfachen Abfragen, zum Beispiel:

data_test <- dbGetQuery(conn, "SELECT * FROM Tweet_Data LIMIT 20")
unique_rows <- dbGetQuery(conn, "SELECT COUNT() AS Total FROM Tweet_Data")
kable(data_test)
SQLite-Datenbank mit den gestreamten Tweets
print(as.numeric(unique_rows))
## [1] 1863


Schritt 3: Analysieren

Wenn das ETL sauber läuft, geht es an Insights und Analyse der gesammelten Daten. Für unsere Tweets probieren wir zwei einfache Dinge: eine Wordcloud der im Tweet-Text genannten Begriffe und eine Zeitreihe, um zu sehen, wann innerhalb der 12 Stunden die meisten Tweets eingingen. Das ist natürlich nur ein kleiner Ausschnitt möglicher Analysen mit Tweets – von Sentimentanalyse bis Psychografie ist vieles denkbar.

Also los: Wir bauen eine hübsche Wordcloud:

# Gather all tweets from the database
all_tweets <- dbGetQuery(conn, "SELECT Tweet_ID, Tweet_Content FROM Tweet_Data")

# Create a term-document matrix and sort the words by frequency
dtm <- TermDocumentMatrix(VCorpus(VectorSource(all_tweets$Tweet_Content)))
dtm_mat <- as.matrix(dtm)
sorted <- sort(rowSums(dtm_mat), decreasing = TRUE)
freq_df <- data.frame(words = names(sorted), freq = sorted)

# Plot the wordcloud
set.seed(42)
wordcloud(words = freq_df$words, freq = freq_df$freq, min.freq = 10,
          max.words=50, random.order=FALSE, rot.per=0.15,
          colors=brewer.pal(8, "RdYlGn"))
Wordcloud der Tweets


Keine große Überraschung: machinelearning und datascience sind die am häufigsten genannten Wörter – es sind schließlich zwei der Hashtags, nach denen wir gestreamt haben. Erwartbar also. Spannender sind die anderen Begriffe. Zum Beispiel tauchen bigdata und artificialintelligence häufig auf, obwohl sie nicht in unseren Keys standen. Man kann also annehmen, dass sie oft zusammen mit den beiden Hauptthemen genannt werden. Auch Wörter wie python oder tensorflow geben zusätzlichen Kontext zu den Inhalten jenseits der Hashtags.

Jetzt eine weitere einfache Analyse: Zu welcher Zeit haben wir in den 12 Stunden die meisten Tweets eingesammelt? Dafür holen wir die Integer-Daten, wandeln sie ins richtige Format und plotten die Tweet-Menge über die Zeit:

# Select the dates in which the tweets were created and convert them into UTC date-time format
all_tweets <- dbGetQuery(conn, "SELECT Tweet_ID, Date_Created FROM Tweet_Data")
all_tweets$Date_Created <- as.POSIXct(all_tweets$Date_Created, origin = "1970-01-01", tz = "UTC")

# Group by the day and hour and count the number of tweets that occurred in each bucket
all_tweets_2 <- all_tweets %>%
    mutate(day = day(Date_Created),
           month = month(Date_Created, label = TRUE),
           hour = hour(Date_Created)) %>%
    mutate(day_hour = paste(month,"-",day,"-",hour, sep = "")) %>%
    group_by(day_hour) %>%
    tally()

# Simple line ggplot
ggplot(all_tweets_2, aes(x = day_hour, y = n)) +
  geom_line(aes(group = 1)) +
  theme_minimal() +
  ggtitle("Tweet Freqeuncy During the 12-h Streming Period")+
  ylab("Tweet Count")+
  xlab("Month-Day-Hour")
Anzahl der Tweets im Zeitverlauf


Super. Wir sehen, dass die meisten einzigartigen Tweets am 2. September zwischen 20:00 und 20:59 Uhr UTC eingingen (20 in 24‑Stunden-Schreibweise).

Fazit

Glückwunsch. Du weißt jetzt, wie du in R eine einfache ETL-Pipeline aufbaust. Die beiden gezeigten Auswertungen sind sehr grundlegende Analysen mit Twitter-Daten. Wie bereits erwähnt, lässt sich noch viel mehr machen, sofern die Pipeline die Daten zuverlässig hineinholt – genau darauf lag hier der Schwerpunkt.

Beachte allerdings: Dieses Tutorial demonstriert nur eine kleine Fallstudie, um den Aufbau von ETL-Pipelines für Twitter-Daten durchzuspielen. Robuste, skalierbare ETL-Pipelines für ein ganzes Unternehmen sind komplex und erfordern umfangreiche Ressourcen und Know-how – besonders, wenn Big Data ins Spiel kommt.

Ich ermutige dich, weiter zu recherchieren und eigene kleine Pipelines zu bauen – gern auch in Python. Vielleicht wagst du dich direkt an Big-Data-Projekte. DataCamp hat zum Beispiel den Kurs Big Data Fundamentals via PySpark, in dem Big Data mit Tools wie PySpark behandelt wird und du dein Wissen vertiefen kannst.

Wenn du mehr über Data Engineering lernen willst, starte mit DataCamps Introduction to Data Engineering. Und wenn du bereit bist, deine neuen Kompetenzen gegenüber Arbeitgebern zu belegen, schau dir unsere Data Engineer Certification an.


Quellen

  1. Foley, D. (2019, 11. Mai). Streaming Twitter Data into a MySQL Database. Abgerufen von https://towardsdatascience.com/streaming-twitter-data-into-a-mysql-database-d62a02b050d6
Themen
R
Datenwissenschaft
Datentechnik
Maschinelles Lernen
SQL

Mehr über R lernen

Kurs

Grundlagen von Big Data mit PySpark

4 Std.
66.5K
Dieser Kurs zeigt praxisnah, wie du in PySpark mit Big Data arbeitest.
Details anzeigenRight Arrow
Kurs Starten
Mehr anzeigenRight Arrow