session_window

Belirli bir zaman damgası sütunu verilen oturum penceresi oluşturur.

Oturum penceresi dinamik pencerelerden biridir; bu da pencerenin uzunluğunun verilen girişlere göre değiştiği anlamına gelir. Oturum penceresinin uzunluğu "oturumun en son girişinin zaman damgası + boşluk süresi" olarak tanımlanır, bu nedenle yeni girişler geçerli oturum penceresine bağlandığında, oturum penceresinin bitiş saati yeni girişlere göre genişletilebilir.

Windows mikrosaniye duyarlığı destekleyebilir. Windows aylar sırasıyla desteklenmez.

Akış sorgusu için, işleme zamanında pencere oluşturmak için işlevini current_timestamp kullanabilirsiniz. gapDuration dize olarak sağlanır, örneğin '1 saniye', '1 gün 12 saat', '2 dakika'. Geçerli aralık dizeleri :'week', 'day', 'hour', 'minute', 'second', 'milisaniye', 'microsecond'.

Giriş satırına göre dinamik olarak boşluk süresi olarak değerlendirilebilen bir Sütun da olabilir.

Çıkış sütunu, 'start' ve 'end' iç içe sütunlarıyla varsayılan olarak 'session_window' adlı bir yapı olacak ve 'start' ve 'end' değerleri olacaktır pyspark.sql.types.TimestampType.

Karşılık gelen Databricks SQL işlevi için bkz session_window . gruplandırma ifadesi.

Sözdizimi

from pyspark.sql import functions as dbf

dbf.session_window(timeColumn=<timeColumn>, gapDuration=<gapDuration>)

Parametreler

Parametre Türü Description
timeColumn pyspark.sql.Column veya str Zamana göre pencereleme için zaman damgası olarak kullanılacak sütun adı veya sütun. Zaman sütunu TimestampType veya TimestampNTZType olmalıdır.
gapDuration pyspark.sql.Column veya literal string Oturumun zaman aşımını belirten python dizesi değişmez değeri veya sütunu. Statik değer, örneğin 10 minutes, 1 secondveya giriş satırına göre aralık süresini dinamik olarak belirten bir ifade/UDF olabilir.

İade

pyspark.sql.Column: hesaplanan sonuçlar için sütun.

Örnekler

from pyspark.sql import functions as dbf
df = spark.createDataFrame([('2016-03-11 09:00:07', 1)], ['dt', 'v'])
df2 = df.groupBy(dbf.session_window('dt', '5 seconds')).agg(dbf.sum('v'))
df2.show(truncate=False)
df2.printSchema()