каналы
Поддержка каналов для потоков.
Примечание: Это часть модуля системы. Не импортируйте его напрямую. Для активации поддержки потоков компилируйте с ключом командной строки --threads:on.
Примечание: Каналы предназначены для типа Thread. Они нестабильны при использовании с spawn
Примечание: Текущая реализация обмена сообщениями не работает с циклическими структурами данных.
Примечание: Каналы нельзя передавать между потоками. Используйте глобальные переменные или передавайте их с помощью ptr.
Пример
Ниже представлен простой пример использования каналов двумя различными способами: блокирующим и неблокирующим.
# Be sure to compile with --threads:on.
# The channels and threads modules are part of system and should not be
# imported.
import os
# Channels can either be:
# - declared at the module level, or
# - passed to procedures by ptr (raw pointer) -- see note on safety.
#
# For simplicity, in this example a channel is declared at module scope.
# Channels are generic, and they include support for passing objects between
# threads.
# Note that objects passed through channels will be deeply copied.
var chan: Channel[string]
# This proc will be run in another thread using the threads module.
proc firstWorker() =
chan.send("Hello World!")
# This is another proc to run in a background thread. This proc takes a while
# to send the message since it sleeps for 2 seconds (or 2000 milliseconds).
proc secondWorker() =
sleep(2000)
chan.send("Another message")
# Initialize the channel.
chan.open()
# Launch the worker.
var worker1: Thread[void]
createThread(worker1, firstWorker)
# Block until the message arrives, then print it out.
echo chan.recv() # "Hello World!"
# Wait for the thread to exit before moving on to the next example.
worker1.joinThread()
# Launch the other worker.
var worker2: Thread[void]
createThread(worker2, secondWorker)
# This time, use a non-blocking approach with tryRecv.
# Since the main thread is not blocked, it could be used to perform other
# useful work while it waits for data to arrive on the channel.
while true:
let tried = chan.tryRecv()
if tried.dataAvailable:
echo tried.msg # "Another message"
break
echo "Pretend I'm doing useful work..."
# For this example, sleep in order not to flood stdout with the above
# message.
sleep(400)
# Wait for the second thread to exit before cleaning up the channel.
worker2.joinThread()
# Clean up the channel.
chan.close() Пример вывода
Программа должна вывести что-то подобное, но имейте в виду, что точные результаты могут отличаться в реальном мире:
Hello World! Pretend I'm doing useful work... Pretend I'm doing useful work... Pretend I'm doing useful work... Pretend I'm doing useful work... Pretend I'm doing useful work... Another message
Безопасная передача каналов
Обратите внимание, что при передаче объектов в процедуры другого потока по указателю (например, через аргумент потока), объекты, созданные с помощью стандартного выделения памяти, будут использовать локальную для потока, управляемую сборщиком мусора, память. Поэтому в целом безопаснее хранить объекты канала в глобальных переменных (как в приведённом выше примере), в этом случае они будут использовать общую для всего процесса (безопасную для потоков) общую кучу.
Однако, можно вручную выделить общую память для каналов, например, с помощью system.allocShared0 и передать эти указатели через аргументы потока:
proc worker(channel: ptr Channel[string]) =
let greeting = channel[].recv()
echo greeting
proc localChannelExample() =
# Use allocShared0 to allocate some shared-heap memory and zero it.
# The usual warnings about dealing with raw pointers apply. Exercise caution.
var channel = cast[ptr Channel[string]](
allocShared0(sizeof(Channel[string]))
)
channel[].open()
# Create a thread which will receive the channel as an argument.
var thread: Thread[ptr Channel[string]]
createThread(thread, worker, channel)
channel[].send("Hello from the main thread!")
# Clean up resources.
thread.joinThread()
channel[].close()
deallocShared(channel)
localChannelExample() # "Hello from the main thread!" Типы
Channel*[TMsg] {...}{.gcsafe.} = RawChannel- канал для межпоточной коммуникации Исходный код Редактировать
Процедуры
proc send*[TMsg](c: var Channel[TMsg]; msg: sink TMsg) {...}{.inline.}- Отправляет сообщение в поток.
msgглубоко копируется. Исходный код Редактировать proc trySend*[TMsg](c: var Channel[TMsg]; msg: sink TMsg): bool {...}{.inline.}-
Пытается отправить сообщение в поток.
msgглубоко копируется. Не блокирует.Возвращает
Исходный код Редактироватьfalseесли сообщение не было отправлено, потому что количество ожидающих элементов в канале превысилоmaxItems. proc recv*[TMsg](c: var Channel[TMsg]): TMsg
-
Получает сообщение из канала
c.Этот процесс блокируется до тех пор, пока не придёт сообщение! Вы можете использовать процедуру peek, чтобы избежать блокировки.
Исходный код Редактировать proc tryRecv*[TMsg](c: var Channel[TMsg]): tuple[dataAvailable: bool, msg: TMsg]
-
Пытается получить сообщение из канала
c, но это может завершиться неудачей по разным причинам, включая конкуренцию.Если это происходит неудачно, возвращает
Исходный код Редактировать(false, default(msg)), иначе возвращает(true, msg). proc peek*[TMsg](c: var Channel[TMsg]): int
-
Возвращает текущее количество сообщений в канале
c.Возвращает -1, если канал закрыт.
Примечание: Это опасно использовать, так как это способствует гонкам. Гораздо лучше использовать tryRecv вместо этого.
Исходный код Редактировать proc open*[TMsg](c: var Channel[TMsg]; maxItems: int = 0)
-
Открывает канал
cдля межпоточной коммуникации.Операция
sendбудет блокироваться до тех пор, пока количество необработанных элементов не станет меньшеmaxItems.Для неограниченной очереди установите
Исходный код РедактироватьmaxItemsв 0. proc close*[TMsg](c: var Channel[TMsg])
- Закрывает канал
cи освобождает связанные с ним ресурсы. Исходный код Редактировать proc ready*[TMsg](c: var Channel[TMsg]): bool
- Возвращает true, если какой-то поток ожидает на канале
cновых сообщений. Исходный код Редактировать
© 2006–2021 Andreas Rumpf
Licensed under the MIT License.
https://nim-lang.org/docs/channels.html